refactor(channels): deslop reply and adapter seams

## What Problem This Solves

Reply orchestration and channel adapters still repeat normalization, routing preflight and send sequencing already owned by shared helpers.

## User Impact

No user-visible behavior change. Reply text, send ordering, receipt reporting, reply targets, private-webchat routing fences and account policies retain their existing behavior.

## Why This Change Was Made

- Use the existing routing decision owner for preflight while preserving lazy runtime loading and the original pre/post-await reads.
- Carry prepared chunk limits into coalescing, use the conversation-label owner directly, and consume LINE's already-normalized media URLs.
- Route SDK text sends through the shared sequence helper, preserving empty chunks, successful undefined results and observer receiver/read timing.
- Remove Signal's single-use attachment formatting wrapper and Zalouser's duplicate sender object; reuse Feishu's local backoff-code owner and the shared receipt normalizers.
- Remove a redundant SDK side-effect import whose value re-export already evaluates the same module.

## Evidence

Independent review completed with no actionable P0–P2 findings after correcting observer forwarding in the shared send loop. Existing assertions are retained; the A2UI fixture readiness fix is described below.

Net **30 production lines removed** across 10 files; one exact shrink-only assertion allowance changes from 3 to 2. One existing UI test gains an event-driven readiness wait; no assertions or cases are removed. No public SDK, configuration, or generated-asset changes.

Blacksmith Testbox proof:
- Candidate `ca6e15f7b39e99e0e4719f4cf2d20c1c58dbe7d7` before the identifier-only lint repair: both cycle checks report **0**; SDK API comparison reports **no changes**; 11 focused files / 320 tests pass; 48 plugin-contract files / 1,132 tests pass; both extension-import inventories report no violations.
- Full sibling coverage on the runtime-identical candidate before reverting private type factoring: 300 files / 4,290 tests pass in 267.24s, with three existing Signal skips. Includes full LINE, Signal, Feishu and Zalouser suites and all dispatch-from-config siblings. The private type factoring was dropped after API comparison detected declaration changes.
- Repaired head `2d1222e5c02406eadc874494e089e4c50a7852bc`: targeted core lint reports zero warnings/errors, 71 SDK reply-payload tests pass, 1,132 plugin-contract tests pass, and both cycle checks report zero. Fresh independent review is clean. This only renames two callback parameters and leaves SDK declarations and runtime behavior unchanged.
- Hosted CI run [36881459169](https://github.com/openclaw/openclaw/actions/runs/36881459169) found two task-caused `eslint(no-shadow)` errors in those callbacks; the repaired head fixes both. The broad changed check independently reached the same error after all type graphs and unused-export scans passed. The repaired-head full `check-changed` replay passed, including all type graphs, all three unused-export scans, lint, and boundary guards. An earlier task-caused assertion ratchet failure was resolved by shrinking the exact Feishu allowance.

The temporary proof checkout was reconstructed from the reviewed Git bundle and verified by exact revision and changed-file hashes after source packaging stalled. No result from a mismatched transport HEAD is counted as proof.

### Fixes found along the way

Hosted CI exposed a browser-fixture race at `board-a2ui.e2e.test.ts:154`: the inline module was inserted but its custom element was not necessarily registered before the first action. Awaiting `customElements.whenDefined` fixes that setup boundary. Production Canvas code and all count/payload assertions are unchanged.

Repair head `fab4fae45be5fd501b7aa21d93ca0ecb90575988` passed all four browser cases and 21 fresh-context executions of the reconnect case on Linux (`--repeats 20 --retry 0`). The full file cost 36.51s wall with one worker; the 21 reconnect executions took 2.115s of test time (29.35s total runner duration). UI E2E types, focused lint, 48 plugin-contract files / 1,132 tests and both zero-cycle checks passed. A negative control restored the original anonymous-listener bug from #161706 in the isolated proof checkout: the unchanged reconnect assertion failed with expected length 1, received 2. The original source was restored and bundles regenerated afterward. The earlier duplicate-config launcher error occurred before tests and was corrected by letting the repository wrapper select the project.

The previous head's remaining stream-reconciliation failure is independently present on hourly main run [36892916454](https://github.com/openclaw/openclaw/actions/runs/36892916454), job `110473226209`: `chat-flow.stream-reconciliation.e2e.test.ts:242`, “reconciles persistence before hydration (tool: false, steer: true)”, omits “The follow-up check is complete.” The mocked fixture authors its own Gateway events and does not execute this PR's changed transport/reply functions. New-head hosted CI will determine the current gate state; no passing replay is claimed as a fix for that separate failure.
This commit is contained in:
Peter Steinberger 2026-10-01 11:23:27 -07:00 • committed by GitHub
parent 1ba5127330
commit 0a2eba4ead
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
12 changed files with 83 additions and 112 deletions

View file

@ -372,7 +372,7 @@ extensions/feishu/src/setup-core.ts 1
extensions/feishu/src/setup-surface.ts 17
extensions/feishu/src/thread-bindings.ts 2
extensions/feishu/src/tool-account.ts 1
extensions/feishu/src/typing-backoff.ts 3
extensions/feishu/src/typing-backoff.ts 2
extensions/feishu/src/wiki.ts 2
extensions/file-transfer/index.ts 4
extensions/file-transfer/src/node-host/dir-fetch.ts 1

View file

@ -21,13 +21,12 @@ export function isFeishuBackoffError(err: unknown): boolean {
if (response.status === 429) {
return true;
}
if (typeof response.data?.code === "number" && FEISHU_BACKOFF_CODES.has(response.data.code)) {
if (getBackoffCodeFromResponse(response.data) !== undefined) {
return true;
}
}
const code = (err as { code?: number }).code;
return typeof code === "number" && FEISHU_BACKOFF_CODES.has(code);
return getBackoffCodeFromResponse(err) !== undefined;
}
export function getBackoffCodeFromResponse(response: unknown): number | undefined {

View file

@ -269,11 +269,7 @@ export async function deliverLineAutoReply(params: {
};
const mediaMessages: messagingApi.Message[] = [];
let deliveryError: unknown;
for (const rawUrl of mediaUrls) {
const url = rawUrl?.trim();
if (!url) {
continue;
}
for (const url of mediaUrls) {
try {
mediaMessages.push(await buildLineMediaMessage(url, mediaOpts, to));
} catch (err) {

View file

@ -4,13 +4,6 @@ import {
} from "openclaw/plugin-sdk/channel-inbound";
import { kindFromMime } from "openclaw/plugin-sdk/media-runtime";
function formatAttachmentKindCount(kind: string, count: number): string {
if (kind === "attachment") {
return `${count} file${count > 1 ? "s" : ""}`;
}
return `${count} ${kind}${count > 1 ? "s" : ""}`;
}
/** Keeps Signal's established multi-attachment text while sharing single-item rendering. */
export function formatSignalMediaText(media: readonly MediaPlaceholderTextFact[]): string {
if (media.length <= 1) {
@ -24,8 +17,8 @@ export function formatSignalMediaText(media: readonly MediaPlaceholderTextFact[]
: (kindFromMime(entry.contentType) ?? "attachment");
kindCounts.set(kind, (kindCounts.get(kind) ?? 0) + 1);
}
const parts = [...kindCounts.entries()].map(([kind, count]) =>
formatAttachmentKindCount(kind, count),
const parts = [...kindCounts.entries()].map(
([kind, count]) => `${count} ${kind === "attachment" ? "file" : kind}${count > 1 ? "s" : ""}`,
);
return `[${parts.join(" + ")} attached]`;
}

View file

@ -174,11 +174,6 @@ const sendZalouserOutbound: NonNullable<ChannelOutboundAdapter["sendMedia"]> = a
}),
);
const zalouserRawSendResultAdapter = {
sendText: sendZalouserOutbound,
sendMedia: sendZalouserOutbound,
};
export const zalouserMessageAdapter = defineChannelMessageAdapter({
id: "zalouser",
durableFinal: {
@ -445,18 +440,15 @@ export const zalouserOutboundAdapter = {
deliveryMode: "direct" as const,
chunker: chunkTextForOutbound,
chunkerMode: "markdown" as const,
sendPayload: async (
ctx: { payload: object } & Parameters<
NonNullable<typeof zalouserRawSendResultAdapter.sendText>
>[0],
) =>
sendPayload: async (ctx: { payload: object } & Parameters<typeof sendZalouserOutbound>[0]) =>
await sendPayloadWithChunkedTextAndMedia({
ctx,
sendText: zalouserRawSendResultAdapter.sendText,
sendMedia: zalouserRawSendResultAdapter.sendMedia,
sendText: sendZalouserOutbound,
sendMedia: sendZalouserOutbound,
emptyResult: createEmptyChannelResult("zalouser"),
}),
...zalouserRawSendResultAdapter,
sendText: sendZalouserOutbound,
sendMedia: sendZalouserOutbound,
sanitizeText: ({ text }) => sanitizeAssistantVisibleText(text),
} satisfies ChannelOutboundAdapter;

View file

@ -100,12 +100,16 @@ export function resolveEffectiveBlockStreamingConfig(params: {
chunking: BlockStreamingChunking;
coalescing: BlockStreamingCoalescing;
} {
const { textLimit } = resolveProviderChunkContext(params.cfg, params.provider, params.accountId);
const providerContext = resolveProviderChunkContext(
params.cfg,
params.provider,
params.accountId,
);
const chunkingDefaults =
params.chunking ?? resolveBlockStreamingChunking(params.cfg, params.provider, params.accountId);
const chunkingMax = clampPositiveInteger(params.maxChunkChars, chunkingDefaults.maxChars, {
min: 1,
max: Math.max(1, textLimit),
max: Math.max(1, providerContext.textLimit),
});
const chunking: BlockStreamingChunking = {
...chunkingDefaults,
@ -114,7 +118,7 @@ export function resolveEffectiveBlockStreamingConfig(params: {
};
const coalescingDefaults = resolveBlockStreamingCoalescing(
params.cfg,
params.provider,
providerContext,
params.accountId,
chunking,
);
@ -164,16 +168,10 @@ export function resolveBlockStreamingChunking(
function resolveBlockStreamingCoalescing(
cfg: OpenClawConfig | undefined,
provider: string | undefined,
{ providerKey, providerId, textLimit }: ReturnType<typeof resolveProviderChunkContext>,
accountId: string | null | undefined,
chunking: BlockStreamingChunking,
): BlockStreamingCoalescing {
const { providerKey, providerId, textLimit } = resolveProviderChunkContext(
cfg,
provider,
accountId,
);
const providerDefaults = providerId
? getChannelPlugin(providerId)?.streaming?.blockStreamingCoalesceDefaults
: undefined;

View file

@ -2,7 +2,6 @@ import { isParentOwnedBackgroundAcpSession } from "@openclaw/acp-core/session-in
import { resolveSendableOutboundReplyParts } from "openclaw/plugin-sdk/reply-payload";
import { readAcpSessionEntryAsync } from "../../acp/runtime/session-meta.js";
import { logVerbose } from "../../globals.js";
import { INTERNAL_MESSAGE_CHANNEL, normalizeMessageChannel } from "../../utils/message-channel.js";
import { resolveCommandTurnTargetSessionKey } from "../command-turn-context.js";
import {
copyReplyPayloadMetadata,
@ -56,25 +55,25 @@ export async function prepareDispatchDelivery(state: GatherDispatchRequestReadyS
? { ...currentAcpSession.entry, acp: currentAcpSession.acp }
: undefined;
const suppressAcpChildUserDelivery = isParentOwnedBackgroundAcpSession(sessionEntryWithAcp);
const normalizedRouteReplyChannel = normalizeMessageChannel(replyRoute.channel);
const normalizedProviderChannel = normalizeMessageChannel(ctx.Provider);
const normalizedSurfaceChannel = normalizeMessageChannel(ctx.Surface);
const normalizedCurrentSurface = normalizedProviderChannel ?? normalizedSurfaceChannel;
const effectiveExplicitDeliverRoute =
ctx.ExplicitDeliverRoute === true || replyRoute.inheritedExternalRoute === true;
const isInternalWebchatTurn =
normalizedCurrentSurface === INTERNAL_MESSAGE_CHANNEL &&
(normalizedSurfaceChannel === INTERNAL_MESSAGE_CHANNEL || !normalizedSurfaceChannel) &&
!effectiveExplicitDeliverRoute;
const hasRouteReplyCandidate = Boolean(
!suppressAcpChildUserDelivery &&
!isInternalWebchatTurn &&
normalizedRouteReplyChannel &&
replyRoute.to &&
normalizedRouteReplyChannel !== normalizedCurrentSurface &&
!state.replyOperationRunState.heartbeat,
);
const routeReplyRuntime = hasRouteReplyCandidate ? await loadRouteReplyRuntime() : undefined;
const {
currentSurface: normalizedCurrentSurface,
isInternalWebchatTurn,
shouldRouteToOriginating: hasRouteReplyCandidate,
} = resolveReplyRoutingDecision({
provider: ctx.Provider,
surface: ctx.Surface,
explicitDeliverRoute: effectiveExplicitDeliverRoute,
originatingChannel: replyRoute.channel,
originatingTo: replyRoute.to,
suppressDirectUserDelivery: suppressAcpChildUserDelivery,
isRoutableChannel: Boolean,
});
const routeReplyRuntime =
hasRouteReplyCandidate && !state.replyOperationRunState.heartbeat
? await loadRouteReplyRuntime()
: undefined;
const {
originatingChannel: routeReplyChannel,
currentSurface,

View file

@ -155,9 +155,7 @@ function finalizeInboundContextImpl<T extends Record<string, unknown>>(
normalized.BodyForAgent = normalized.agentText;
normalized.BodyForCommands = normalized.commandText;
const label =
normalizeOptionalString(normalized.ConversationLabel) ??
normalizeOptionalString(resolveConversationLabel(normalized));
const label = resolveConversationLabel(normalized);
if (label) {
normalized.ConversationLabel = label;
}

View file

@ -2,7 +2,7 @@ import { asOptionalRecord } from "@openclaw/normalization-core/record-coerce";
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
import {
normalizeUniqueStringEntries,
uniqueStrings,
normalizeUniqueTrimmedStringList,
} from "@openclaw/normalization-core/string-normalization";
import type {
MessageReceipt,
@ -14,9 +14,6 @@ type MessageReceiptInputResult = MessageReceiptSourceResult & {
receipt?: MessageReceipt;
};
const normalizeIdentity = (value: string | undefined): string | undefined =>
value?.trim() || undefined;
/** Reads reported recipients, including every physical part of an aggregate receipt. */
export function listMessageReceiptSourceTargets(value: unknown): string[] {
const targets = new Set<string>();
@ -64,9 +61,9 @@ export function resolveReceiptSourceId(result: MessageReceiptInputResult): strin
return undefined;
}
return (
normalizeIdentity(result.messageId) ??
normalizeOptionalString(result.messageId) ??
(result.receipt ? resolveMessageReceiptPrimaryId(result.receipt) : undefined) ??
normalizeIdentity(result.pollId)
normalizeOptionalString(result.pollId)
);
}
@ -79,21 +76,24 @@ export function createMessageReceiptFromOutboundResults(params: {
sentAt?: number;
}): MessageReceipt {
const sentResults = params.results.filter((result) => result.outcome !== "not_sent");
const requestedThreadId = normalizeIdentity(params.threadId);
const requestedThreadId = normalizeOptionalString(params.threadId);
const providerThreadIds = normalizeUniqueStringEntries(
sentResults.flatMap(({ receipt }) =>
receipt?.parts.length
? receipt.parts.flatMap(
(part) => normalizeIdentity(part.threadId) ?? normalizeIdentity(receipt.threadId) ?? [],
(part) =>
normalizeOptionalString(part.threadId) ??
normalizeOptionalString(receipt.threadId) ??
[],
)
: (normalizeIdentity(receipt?.threadId) ?? []),
: (normalizeOptionalString(receipt?.threadId) ?? []),
),
);
const aggregateThreadId =
providerThreadIds.length > 1 ? undefined : (providerThreadIds[0] ?? requestedThreadId);
const parts = sentResults.flatMap((result, resultIndex) => {
if (result.receipt) {
const receiptThreadId = normalizeIdentity(result.receipt.threadId) ?? requestedThreadId;
const receiptThreadId = normalizeOptionalString(result.receipt.threadId) ?? requestedThreadId;
if (result.receipt.parts.length === 0) {
return result.receipt.platformMessageIds.map((platformMessageId, partIndex) => ({
platformMessageId,
@ -109,7 +109,7 @@ export function createMessageReceiptFromOutboundResults(params: {
return result.receipt.parts.map((part, partIndex) => ({
...part,
index: part.index ?? partIndex,
...(normalizeIdentity(part.threadId) || !receiptThreadId
...(normalizeOptionalString(part.threadId) || !receiptThreadId
? {}
: { threadId: receiptThreadId }),
...(part.replyToId || !params.replyToId || hasPartReplyMetadata
@ -132,19 +132,16 @@ export function createMessageReceiptFromOutboundResults(params: {
},
];
});
const platformMessageIds = uniqueStrings(
sentResults
.flatMap((result) =>
result.receipt
? [
result.receipt.primaryPlatformMessageId,
...result.receipt.platformMessageIds,
...result.receipt.parts.map((part) => part.platformMessageId),
]
: [resolveReceiptSourceId(result)],
)
.map(normalizeIdentity)
.filter((id): id is string => Boolean(id)),
const platformMessageIds = normalizeUniqueTrimmedStringList(
sentResults.flatMap((result) =>
result.receipt
? [
result.receipt.primaryPlatformMessageId,
...result.receipt.platformMessageIds,
...result.receipt.parts.map((part) => part.platformMessageId),
]
: [resolveReceiptSourceId(result)],
),
);
const firstNestedReceipt = sentResults.find((result) => result.receipt)?.receipt;
return {
@ -167,13 +164,13 @@ export function listMessageReceiptPlatformIds(receipt: MessageReceipt): string[]
/** Resolves the explicit primary platform id, falling back to the first unique receipt id. */
export function resolveMessageReceiptPrimaryId(receipt: MessageReceipt): string | undefined {
const primary = normalizeIdentity(receipt.primaryPlatformMessageId);
const primary = normalizeOptionalString(receipt.primaryPlatformMessageId);
if (primary) {
return primary;
}
return (
listMessageReceiptPlatformIds(receipt)[0] ??
receipt.parts.map((part) => normalizeIdentity(part.platformMessageId)).find(Boolean)
receipt.parts.map((part) => normalizeOptionalString(part.platformMessageId)).find(Boolean)
);
}
@ -183,12 +180,14 @@ export function resolveMessageReceiptThreadId(
requestedThreadId?: string,
): string | undefined {
const partThreadIds = normalizeUniqueStringEntries(
receipt.parts.flatMap((part) => normalizeIdentity(part.threadId) ?? []),
receipt.parts.flatMap((part) => normalizeOptionalString(part.threadId) ?? []),
);
if (partThreadIds.length > 1) {
return undefined;
}
return (
partThreadIds[0] ?? normalizeIdentity(receipt.threadId) ?? normalizeIdentity(requestedThreadId)
partThreadIds[0] ??
normalizeOptionalString(receipt.threadId) ??
normalizeOptionalString(requestedThreadId)
);
}

View file

@ -15,7 +15,6 @@ import {
type FinalizeChannelInboundContextResult,
} from "../channels/inbound-event/context.js";
import type { InboundEventKind } from "../channels/inbound-event/kind.js";
import "../channels/turn/dispatch-result.js";
import { runPreparedChannelTurn } from "../channels/turn/execution.js";
import {
dispatchAssembledChannelTurn,

View file

@ -262,13 +262,11 @@ export async function sendPayloadWithChunkedTextAndMedia<
}
const limit = params.textChunkLimit;
const chunkedText = limit && params.chunker ? params.chunker(text, limit) : [text];
const chunks = resolveTextChunksWithFallback(text, chunkedText);
let lastResult = params.emptyResult;
for (const chunk of chunks) {
lastResult = await params.sendText({ ...params.ctx, text: chunk });
await params.onResult?.(lastResult);
}
return lastResult;
return (await sendPayloadTextChunkSequence({
chunks: resolveTextChunksWithFallback(text, chunkedText),
send: ({ text: chunk }) => params.sendText({ ...params.ctx, text: chunk }),
onResult: (result) => params.onResult?.(result),
}))!;
}
/**
@ -455,19 +453,18 @@ export async function sendTextMediaPayload(params: {
limit && params.adapter.chunker
? params.adapter.chunker(text, limit, { formatting: params.ctx.formatting })
: [text];
const chunks = resolveTextChunksWithFallback(text, chunkedText);
let lastResult: Awaited<ReturnType<NonNullable<typeof params.adapter.sendText>>>;
for (const chunk of chunks) {
lastResult = await sendAndReport((onDeliveryResult) =>
params.adapter.sendText!({
...params.ctx,
text: chunk,
replyToId: nextReplyToId(),
onDeliveryResult,
}),
);
}
return lastResult!;
return (await sendPayloadTextChunkSequence({
chunks: resolveTextChunksWithFallback(text, chunkedText),
send: ({ text: chunk }) =>
sendAndReport((onDeliveryResult) =>
params.adapter.sendText!({
...params.ctx,
text: chunk,
replyToId: nextReplyToId(),
onDeliveryResult,
}),
),
}))!;
}
/** Detect numeric-looking target ids for channels that distinguish ids from handles. */

View file

@ -120,7 +120,8 @@ describeControlUiE2e("Control UI dashboard A2UI", () => {
path: path.resolve("extensions/canvas/src/host/a2ui/a2ui.bundle.js"),
type: "module",
});
const result = await page.evaluate(() => {
const result = await page.evaluate(async () => {
await customElements.whenDefined("openclaw-a2ui-host");
const emitted: unknown[] = [];
Reflect.set(globalThis, "openclaw", {
state: {