fix(telegram): keep tool progress transient without drafts

Keep quiet Telegram tool progress on the transient draft lane when it exists. Preserve durable delivery for media, execution approvals, and ask-user prompts across direct and queued turns.

Co-authored-by: Ayaan Zaidi <hi@obviy.us>
This commit is contained in:
Ayaan Zaidi
2026-08-09 19:18:02 +05:30
committed by GitHub
parent 4fec524243
commit 6a586cd0ea
9 changed files with 224 additions and 78 deletions
@@ -249,6 +249,7 @@ export async function runTelegramDispatchTurn(params: {
suppressDefaultToolProgressMessages:
!params.draft.streamDeliveryEnabled || Boolean(params.draft.answerLane.stream),
forceToolResultProgress:
Boolean(params.draft.answerLane.stream) &&
params.streamMode === "progress" &&
resolveChannelStreamingPreviewToolProgress(
params.telegramCfg,
@@ -46,6 +46,7 @@ describeTelegramDispatch("Telegram provider preview hook safety", () => {
replyOptions: expect.objectContaining({
onPartialReply: undefined,
disableBlockStreaming: undefined,
forceToolResultProgress: false,
}),
});
},
@@ -28,7 +28,10 @@ import {
runWithDispatchAbortSignal,
} from "./dispatch-from-config.abort.js";
import { createReplyDispatchEvent } from "./dispatch-from-config.events.js";
import { readAskUserQuestionId } from "./dispatch-from-config.payloads.js";
import {
hasExecApprovalPayload,
requiresDurableToolResultDelivery,
} from "./dispatch-from-config.payloads.js";
import { extendPreparedDispatchState } from "./dispatch-from-config.phase-state.js";
import type { PrepareDispatchOperationReadyState } from "./dispatch-from-config.prepare-operation.js";
import {
@@ -133,25 +136,11 @@ export async function chooseDispatchRoute(state: PrepareDispatchOperationReadySt
ctx.InboundEventKind !== "room_event" &&
!sendPolicyDenied;
let finalReplyDeliveryStarted = false;
const hasExecApprovalPayload = (payload: ReplyPayload) => {
const execApproval =
payload.channelData &&
typeof payload.channelData === "object" &&
!Array.isArray(payload.channelData)
? payload.channelData.execApproval
: undefined;
return execApproval && typeof execApproval === "object" && !Array.isArray(execApproval);
};
const hasAskUserPayload = (payload: ReplyPayload) => {
const askUser = payload.channelData?.askUser;
return askUser && typeof askUser === "object" && !Array.isArray(askUser);
};
const shouldSuppressLateTextOnlyToolProgress = (payload: ReplyPayload) => {
if (!finalReplyDeliveryStarted) {
return false;
}
const reply = resolveSendableOutboundReplyParts(payload);
return !reply.hasMedia && !hasExecApprovalPayload(payload) && !hasAskUserPayload(payload);
return !requiresDurableToolResultDelivery(payload);
};
// Durable inter-tool commentary lane: with verbose progress on, preamble
// items become standalone progress messages like tool summaries. The latest
@@ -694,9 +683,6 @@ export async function chooseDispatchRoute(state: PrepareDispatchOperationReadySt
shouldDeliverVerboseProgressDespiteSourceSuppression,
shouldDeliverForcedToolProgressDespiteSourceSuppression,
shouldDeliverFastModeAutoProgressDespiteSourceSuppression,
hasExecApprovalPayload,
hasAskUserPayload,
readAskUserQuestionId,
shouldSuppressLateTextOnlyToolProgress,
flushPendingCommentaryProgress,
noteCommentaryProgress,
@@ -1,7 +1,6 @@
import {
hasOutboundReplyContent,
isFastModeAutoProgressPayload,
resolveSendableOutboundReplyParts,
} from "openclaw/plugin-sdk/reply-payload";
import { isAskUserPromptPending } from "../../agents/tools/ask-user-tool.js";
import { normalizeAgentPlanSteps } from "../../channels/streaming.js";
@@ -22,7 +21,12 @@ import {
type InternalReplyResolverOptions,
createReplyDispatchEvent,
} from "./dispatch-from-config.events.js";
import { shouldDeliverDespiteSourceReplySuppression } from "./dispatch-from-config.payloads.js";
import {
hasAskUserPayload,
readAskUserQuestionId,
requiresDurableToolResultDelivery,
shouldDeliverDespiteSourceReplySuppression,
} from "./dispatch-from-config.payloads.js";
import { extendPreparedDispatchState } from "./dispatch-from-config.phase-state.js";
import type { PrepareDispatchExecutionReadyState } from "./dispatch-from-config.prepare-execution.js";
import { waitForReplyDispatcherIdle } from "./reply-dispatcher.js";
@@ -41,7 +45,6 @@ export async function executeDispatch(state: PrepareDispatchExecutionReadyState)
flushPendingCommentaryProgress,
getDispatchAbortOperation,
getDispatchAbortSignal,
hasAskUserPayload,
hookRunner,
isDispatchOperationAborted,
markInboundDedupeReplayUnsafe,
@@ -216,20 +219,33 @@ export async function executeDispatch(state: PrepareDispatchExecutionReadyState)
state.shouldDeliverFastModeAutoProgressDespiteSourceSuppression();
const isForcedToolProgress =
state.shouldDeliverForcedToolProgressDespiteSourceSuppression();
const progressCallbackForwarded = state.shouldForwardToolResultProgressCallback(
payload,
isFastModeAutoProgress,
);
if (progressCallbackForwarded) {
await onToolResultFromReplyOptions?.(payload);
const forceToolResultProgress =
params.replyOptions?.forceToolResultProgress === true;
const requiresDurableToolResult =
forceToolResultProgress && requiresDurableToolResultDelivery(payload);
const shouldForwardToolResultProgress = isFastModeAutoProgress
? shouldForwardProgressCallback({
forwardWhenSourceDeliverySuppressed: true,
})
: forceToolResultProgress
? !requiresDurableToolResult &&
!state.shouldEmitVerboseProgress() &&
shouldForwardProgressCallback({
forwardWhenSourceDeliverySuppressed: true,
})
: state.shouldSendToolSummaries() && shouldForwardProgressCallback();
const toolResultProgressCallback = shouldForwardToolResultProgress
? onToolResultFromReplyOptions
: undefined;
if (toolResultProgressCallback) {
await toolResultProgressCallback(payload);
}
if (isDispatchOperationAborted()) {
return;
}
if (
isFastModeAutoProgress &&
progressCallbackForwarded &&
onToolResultFromReplyOptions
toolResultProgressCallback &&
(isFastModeAutoProgress || forceToolResultProgress)
) {
return;
}
@@ -284,19 +300,14 @@ export async function executeDispatch(state: PrepareDispatchExecutionReadyState)
!isFastModeAutoProgressPayload(deliveryPayload) &&
!isForcedToolProgress
) {
const hasMedia = resolveSendableOutboundReplyParts(deliveryPayload).hasMedia;
if (
!hasMedia &&
!state.hasExecApprovalPayload(deliveryPayload) &&
!hasAskUserPayload(deliveryPayload)
) {
if (!requiresDurableToolResultDelivery(deliveryPayload)) {
return;
}
}
if (deliveryPayload.isError === true) {
markVisibleToolErrorProgress();
}
const askUserQuestionId = state.readAskUserQuestionId(deliveryPayload);
const askUserQuestionId = readAskUserQuestionId(deliveryPayload);
if (
askUserQuestionId !== undefined &&
!(await isAskUserPromptPending(askUserQuestionId))
@@ -2608,7 +2608,7 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () =>
expect(dispatcher.sendFinalReply).not.toHaveBeenCalled();
});
it("forwards suppressed tool progress callbacks in message-tool-only mode", async () => {
it("lets the channel own forced tool progress at verbosity off", async () => {
setNoAbort();
sessionStoreMocks.currentEntry = {
sessionId: "s1",
@@ -2616,9 +2616,10 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () =>
sendPolicy: "allow",
};
const dispatcher = createDispatcher();
const onToolResult = vi.fn();
const onToolResult = vi.fn(() => false);
const payload = { text: "🧠 Memory Search: release notes" } satisfies ReplyPayload;
const replyResolver = vi.fn(async (_ctx: MsgContext, opts?: GetReplyOptions) => {
await opts?.onToolResult?.({ text: "🛠️ Exec: ruby sleep proof" });
await opts?.onToolResult?.(payload);
return { text: "NO_REPLY" } satisfies ReplyPayload;
});
const ctx = buildTestCtx({ SessionKey: "test:session", ChatType: "channel" });
@@ -2630,6 +2631,7 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () =>
replyResolver,
replyOptions: {
sourceReplyDeliveryMode: "message_tool_only",
forceToolResultProgress: true,
suppressDefaultToolProgressMessages: true,
allowProgressCallbacksWhenSourceDeliverySuppressed: true,
onToolResult,
@@ -2638,11 +2640,80 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () =>
expect(result.queuedFinal).toBe(false);
expect(result.sourceReplyDeliveryMode).toBe("message_tool_only");
expect(onToolResult).toHaveBeenCalledWith({ text: "🛠️ Exec: ruby sleep proof" });
expect(onToolResult).toHaveBeenCalledOnce();
expect(onToolResult).toHaveBeenCalledWith(payload);
expect(dispatcher.sendToolResult).not.toHaveBeenCalled();
expect(dispatcher.sendFinalReply).not.toHaveBeenCalled();
});
it.each([
{
label: "media",
payload: { mediaUrl: "https://example.com/tool-result.png" },
},
{
label: "captioned media",
payload: {
text: "Generated image",
mediaUrl: "https://example.com/tool-result.png",
},
},
{
label: "exec approvals",
payload: {
text: "Approval required.",
channelData: {
execApproval: {
approvalId: "117ba06d-1111-2222-3333-444444444444",
approvalSlug: "117ba06d",
allowedDecisions: ["allow-once", "allow-always", "deny"],
},
},
},
},
{
label: "ask-user prompts",
payload: {
text: "Question for you: Where should this deploy?",
channelData: { askUser: { questionId: "question-owned-by-agent-runtime" } },
},
},
] satisfies Array<{ label: string; payload: ReplyPayload }>)(
"keeps forced $label durable when channel progress is available",
async ({ payload }) => {
setNoAbort();
sessionStoreMocks.currentEntry = {
sessionId: "s1",
updatedAt: 0,
sendPolicy: "allow",
};
const dispatcher = createDispatcher();
const onToolResult = vi.fn(() => false);
const replyResolver = vi.fn(async (_ctx: MsgContext, opts?: GetReplyOptions) => {
await opts?.onToolResult?.(payload);
return { text: "NO_REPLY" } satisfies ReplyPayload;
});
await dispatchReplyFromConfig({
ctx: buildTestCtx({ SessionKey: "test:session", ChatType: "channel" }),
cfg: emptyConfig,
dispatcher,
replyResolver,
replyOptions: {
sourceReplyDeliveryMode: "message_tool_only",
forceToolResultProgress: true,
suppressDefaultToolProgressMessages: true,
allowProgressCallbacksWhenSourceDeliverySuppressed: true,
onToolResult,
},
});
expect(onToolResult).not.toHaveBeenCalled();
expect(dispatcher.sendToolResult).toHaveBeenCalledWith(payload);
expect(dispatcher.sendFinalReply).not.toHaveBeenCalled();
},
);
it("delivers forced tool progress in message-tool-only mode without verbose progress", async () => {
setNoAbort();
sessionStoreMocks.currentEntry = {
@@ -2685,6 +2756,7 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () =>
verboseLevel: "on",
};
const dispatcher = createDispatcher();
const onToolResult = vi.fn(() => false);
const replyResolver = vi.fn(async (_ctx: MsgContext, opts?: GetReplyOptions) => {
await opts?.onToolResult?.({ text: "🛠️ Exec: echo post-restart" });
return { text: "NO_REPLY" } satisfies ReplyPayload;
@@ -2698,6 +2770,9 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () =>
replyResolver,
replyOptions: {
sourceReplyDeliveryMode: "message_tool_only",
forceToolResultProgress: true,
allowProgressCallbacksWhenSourceDeliverySuppressed: true,
onToolResult,
},
});
@@ -2706,6 +2781,7 @@ describe("sendPolicy deny — suppress delivery, not processing (#53328)", () =>
expect(dispatcher.sendToolResult).toHaveBeenCalledWith(
expect.objectContaining({ text: "🛠️ Exec: echo post-restart" }),
);
expect(onToolResult).not.toHaveBeenCalled();
expect(dispatcher.sendFinalReply).not.toHaveBeenCalled();
});
@@ -1,3 +1,4 @@
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
import { resolveSendableOutboundReplyParts } from "openclaw/plugin-sdk/reply-payload";
@@ -43,13 +44,29 @@ export function shouldDeliverDespiteSourceReplySuppression(
export function readAskUserQuestionId(payload: ReplyPayload): string | undefined {
const askUser = payload.channelData?.askUser;
if (!askUser || typeof askUser !== "object" || Array.isArray(askUser)) {
if (!isRecord(askUser)) {
return undefined;
}
const questionId = (askUser as { questionId?: unknown }).questionId;
const questionId = askUser.questionId;
return typeof questionId === "string" ? questionId : undefined;
}
export function hasExecApprovalPayload(payload: ReplyPayload): boolean {
return isRecord(payload.channelData?.execApproval);
}
export function hasAskUserPayload(payload: ReplyPayload): boolean {
return isRecord(payload.channelData?.askUser);
}
export function requiresDurableToolResultDelivery(payload: ReplyPayload): boolean {
return (
resolveSendableOutboundReplyParts(payload).hasMedia ||
hasExecApprovalPayload(payload) ||
hasAskUserPayload(payload)
);
}
export function createFinalDispatchPayloadDedupeKey(payload: ReplyPayload): string {
const metadata = getReplyPayloadMetadata(payload);
return JSON.stringify({
@@ -1,4 +1,3 @@
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
import {
isFastModeAutoProgressPayload,
resolveSendableOutboundReplyParts,
@@ -15,6 +14,7 @@ import { normalizeMessageChannel } from "../../utils/message-channel.js";
import type { GetReplyOptions } from "../get-reply-options.types.js";
import type { ReplyPayload } from "../reply-payload.js";
import type { ChooseDispatchRouteReadyState } from "./dispatch-from-config.choose-route.js";
import { hasAskUserPayload, hasExecApprovalPayload } from "./dispatch-from-config.payloads.js";
import { extendPreparedDispatchState } from "./dispatch-from-config.phase-state.js";
import { loadGetReplyFromConfigRuntime } from "./dispatch-from-config.runtime-loaders.js";
import {
@@ -126,16 +126,10 @@ export async function prepareDispatchExecution(state: ChooseDispatchRouteReadySt
if (shouldSendToolSummaries()) {
return payload;
}
const execApproval =
payload.channelData &&
typeof payload.channelData === "object" &&
!Array.isArray(payload.channelData)
? payload.channelData.execApproval
: undefined;
if (execApproval && typeof execApproval === "object" && !Array.isArray(execApproval)) {
if (hasExecApprovalPayload(payload)) {
return payload;
}
if (state.hasAskUserPayload(payload)) {
if (hasAskUserPayload(payload)) {
return payload;
}
if (isFastModeAutoProgressPayload(payload)) {
@@ -197,25 +191,6 @@ export async function prepareDispatchExecution(state: ChooseDispatchRouteReadySt
const onPatchSummaryFromReplyOptions = params.replyOptions?.onPatchSummary;
const allowSuppressedSourceProgressCallbacks =
params.replyOptions?.allowProgressCallbacksWhenSourceDeliverySuppressed === true;
const isChannelOwnedToolResultProgressPayload = (payload: ReplyPayload) => {
const text = normalizeOptionalString(payload.text);
return Boolean(text?.startsWith("🛠️") || text?.startsWith("🔧"));
};
const shouldForwardToolResultProgressCallback = (
payload: ReplyPayload,
isFastModeAutoProgress: boolean,
) => {
if (isFastModeAutoProgress) {
return shouldForwardProgressCallback({ forwardWhenSourceDeliverySuppressed: true });
}
if (
allowSuppressedSourceProgressCallbacks &&
isChannelOwnedToolResultProgressPayload(payload)
) {
return shouldForwardProgressCallback({ forwardWhenSourceDeliverySuppressed: true });
}
return shouldSendToolSummaries() && shouldForwardProgressCallback();
};
const shouldAllowQuietChannelOwnedProgressCallbacks = (options?: {
allowWhenToolSummariesHidden?: boolean;
requiresToolSummaryVisibility?: boolean;
@@ -433,7 +408,6 @@ export async function prepareDispatchExecution(state: ChooseDispatchRouteReadySt
onPlanUpdateFromReplyOptions,
onApprovalEventFromReplyOptions,
onPatchSummaryFromReplyOptions,
shouldForwardToolResultProgressCallback,
waitForPendingDirectBlockReplyDelivery,
shouldForwardProgressCallback,
preserveProgressCallbackStartOrder,
@@ -1,4 +1,5 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import type { ReplyPayload } from "../types.js";
import type { AgentTurnParams } from "./agent-runner-execution.types.js";
import type { AdmittedFollowupTurn } from "./followup-turn-admission.js";
@@ -303,6 +304,79 @@ describe("executeFollowupTurn", () => {
expect(onDurableToolResult).not.toHaveBeenCalled();
});
it.each([
{
label: "media",
payload: { mediaUrl: "https://example.com/tool-result.png" },
},
{
label: "captioned media",
payload: {
text: "Generated image",
mediaUrl: "https://example.com/tool-result.png",
},
},
{
label: "exec approvals",
payload: {
text: "Approval required.",
channelData: {
execApproval: {
approvalId: "117ba06d-1111-2222-3333-444444444444",
approvalSlug: "117ba06d",
allowedDecisions: ["allow-once", "allow-always", "deny"],
},
},
},
},
{
label: "ask-user prompts",
payload: {
text: "Question for you: Where should this deploy?",
channelData: { askUser: { questionId: "question-owned-by-agent-runtime" } },
},
},
] satisfies Array<{ label: string; payload: ReplyPayload }>)(
"keeps quiet forced $label on the durable path",
async ({ payload }) => {
const onChannelToolResult = vi.fn(async () => {});
const onDurableToolResult = vi.fn(async () => {});
const turn = createTurn({
session: {
kind: "session",
key: "main",
current: () => ({ sessionId: "session", updatedAt: 1, verboseLevel: "off" }),
publish: () => undefined,
adopt: () => undefined,
},
});
state.execute.mockImplementation(async (params: AgentTurnParams) => {
await params.opts?.onToolResult?.(payload);
return { runId: "run-1", outcome: { kind: "rejected", payload: { text: "done" } } };
});
const result = await executeFollowupTurn({
turn,
defaults: {
typing: createTypingController(),
typingMode: "never",
defaultModel: "claude",
opts: {
forceToolResultProgress: true,
onToolResult: onChannelToolResult,
},
},
onToolResult: onDurableToolResult,
onCompactionNoticePayload: vi.fn(async () => {}),
});
await result.progress.drain();
expect(onChannelToolResult).not.toHaveBeenCalled();
expect(onDurableToolResult).toHaveBeenCalledOnce();
expect(onDurableToolResult).toHaveBeenCalledWith(payload, { runId: "run-1" });
},
);
it("keeps verbose tool results durable when channel progress is available", async () => {
const onChannelToolResult = vi.fn(async () => {});
const onDurableToolResult = vi.fn(async () => {});
@@ -7,6 +7,7 @@ import type { ReplyPayload } from "../types.js";
import { executeAgentTurn } from "./agent-runner-execution.js";
import type { AgentTurnExecutionResult } from "./agent-runner-execution.types.js";
import { resetReplyRunSession } from "./agent-runner-session-reset.js";
import { requiresDurableToolResultDelivery } from "./dispatch-from-config.payloads.js";
import type { AdmittedFollowupTurn, FollowupRunnerParams } from "./followup-turn-admission.js";
import type { InternalGetReplyOptions } from "./get-reply.types.js";
import { createTypingSignaler, type TypingSignaler } from "./typing-mode.js";
@@ -275,17 +276,22 @@ export async function executeFollowupTurn(params: {
if (!progressAllowed()) {
return false;
}
const verboseToolResult = shouldEmitVerboseToolResult();
const toolResultProgressVisible = Boolean(channelToolResultProgress) || verboseToolResult;
const requiresDurableToolResult = requiresDurableToolResultDelivery(payload);
const verboseToolResult = !requiresDurableToolResult && shouldEmitVerboseToolResult();
const transientToolResultProgress = requiresDurableToolResult
? undefined
: channelToolResultProgress;
const toolResultDeliveryAvailable =
Boolean(transientToolResultProgress) || verboseToolResult || requiresDurableToolResult;
if (
turn.queued.run.sourceReplyDeliveryMode === "message_tool_only" &&
!toolResultProgressVisible
!toolResultDeliveryAvailable
) {
return false;
}
const visible =
channelToolResultProgress && !verboseToolResult
? (await settleProgressVisibilityCallbackResult(channelToolResultProgress(payload)))
transientToolResultProgress && !verboseToolResult
? (await settleProgressVisibilityCallbackResult(transientToolResultProgress(payload)))
.visible
: await params.onToolResult(payload, { runId: turn.runId }).then(() => true);
if (visible && payload.isError === true) {