fix(dev): clean Telegram flow previews on failure

This commit is contained in:
Vincent Koc
2026-06-01 17:37:15 +02:00
parent 9cb347e4c3
commit 461999c060
2 changed files with 121 additions and 15 deletions
+59 -15
View File
@@ -51,6 +51,10 @@ type TelegramFlowResult = {
previewUpdates: number;
};
function toError(value: unknown): Error {
return value instanceof Error ? value : new Error(String(value));
}
type TelegramThinkingFinalDeps = {
createDraftStream?: (params: {
accountId?: string;
@@ -310,15 +314,35 @@ export async function runTelegramThinkingFinalFlow(
});
const wait = deps.sleep ?? sleep;
for (const update of thinkingUpdates) {
stream.update(formatReasoningMessage(update));
await stream.flush();
if (delayMs > 0) {
await wait(delayMs);
let previewStarted = false;
let flowError: unknown;
try {
for (const update of thinkingUpdates) {
previewStarted = true;
stream.update(formatReasoningMessage(update));
await stream.flush();
if (delayMs > 0) {
await wait(delayMs);
}
}
} catch (error) {
flowError = error;
}
let cleanupError: unknown;
if (previewStarted) {
try {
await stream.clear();
} catch (error) {
cleanupError = error;
}
}
if (flowError) {
throw toError(flowError);
}
if (cleanupError) {
throw toError(cleanupError);
}
await stream.clear();
const final = await (deps.sendFinal ?? sendTelegramFinal)({
accountId: options.accountId,
cfg: options.cfg,
@@ -350,19 +374,39 @@ export async function runTelegramWorkingFinalFlow(
let previewUpdates = 0;
let lastPreviewText = "";
const updateIntervalMs = delayMs > 0 ? delayMs : 1_000;
for (let elapsedMs = 0; elapsedMs < durationMs; elapsedMs += updateIntervalMs) {
const previewText = formatWorkingProgressPreview(elapsedMs);
if (previewText !== lastPreviewText) {
await draft.update(previewText);
lastPreviewText = previewText;
previewUpdates += 1;
let draftStarted = false;
let flowError: unknown;
try {
for (let elapsedMs = 0; elapsedMs < durationMs; elapsedMs += updateIntervalMs) {
const previewText = formatWorkingProgressPreview(elapsedMs);
if (previewText !== lastPreviewText) {
draftStarted = true;
await draft.update(previewText);
lastPreviewText = previewText;
previewUpdates += 1;
}
if (delayMs > 0 && elapsedMs + updateIntervalMs < durationMs) {
await wait(delayMs);
}
}
if (delayMs > 0 && elapsedMs + updateIntervalMs < durationMs) {
await wait(delayMs);
} catch (error) {
flowError = error;
}
let cleanupError: unknown;
if (draftStarted) {
try {
draft.stop();
} catch (error) {
cleanupError = error;
}
}
if (flowError) {
throw toError(flowError);
}
if (cleanupError) {
throw toError(cleanupError);
}
draft.stop();
const final = await (deps.sendFinal ?? sendTelegramFinal)({
accountId: options.accountId,
cfg: options.cfg,
@@ -107,6 +107,39 @@ describe("channel message flows dev runner", () => {
expect(result).toEqual({ finalMessageId: "99", previewUpdates: 3 });
});
it("clears thinking previews when streaming fails before the final answer", async () => {
const stream = {
update: vi.fn(() => {}),
flush: vi.fn(async () => {
throw new Error("flush failed");
}),
clear: vi.fn(async () => {}),
stop: vi.fn(async () => {}),
messageId: vi.fn(() => 17),
forceNewMessage: vi.fn(),
};
const sendFinal = vi.fn(async () => ({ messageId: "99", chatId: "123" }));
await expect(
runTelegramThinkingFinalFlow(
{
cfg: {} as OpenClawConfig,
delayMs: 0,
target: "123",
thinkingUpdates: ["Checking the request."],
},
{
createDraftStream: vi.fn(() => stream),
sendFinal,
sleep: vi.fn(async () => {}),
},
),
).rejects.toThrow("flush failed");
expect(stream.clear).toHaveBeenCalledOnce();
expect(sendFinal).not.toHaveBeenCalled();
});
it("streams working updates through native message drafts before the final answer", async () => {
const draft = {
update: vi.fn(async () => true),
@@ -151,6 +184,35 @@ describe("channel message flows dev runner", () => {
expect(result).toEqual({ finalMessageId: "100", previewUpdates: 6 });
});
it("stops native working drafts when progress updates fail before the final answer", async () => {
const draft = {
update: vi.fn(async () => {
throw new Error("draft update failed");
}),
stop: vi.fn(async () => {}),
};
const sendFinal = vi.fn(async () => ({ messageId: "100", chatId: "123" }));
await expect(
runTelegramWorkingFinalFlow(
{
cfg: {} as OpenClawConfig,
delayMs: 0,
durationMs: 12_000,
target: "123",
},
{
createNativeToolProgressDraft: vi.fn(() => draft),
sendFinal,
sleep: vi.fn(async () => {}),
},
),
).rejects.toThrow("draft update failed");
expect(draft.stop).toHaveBeenCalledOnce();
expect(sendFinal).not.toHaveBeenCalled();
});
it("uses two second progress update cadence by default", async () => {
const draft = {
update: vi.fn(async () => true),