fix(telegram): reset delivery owner for successive answers

This commit is contained in:
Ayaan Zaidi 2026-09-23 01:06:13 +05:30
parent 690d2f8ce5
commit da8efdbc09
No known key found for this signature in database
3 changed files with 90 additions and 12 deletions

View file

@ -244,7 +244,7 @@ eligible for normal fallback.
| `observeSuppression()` | Record an intentional final no-send decision, such as cancellation by an outbound modifier hook. It cannot hide an existing delivery failure or revoke accepted content. |
| `cleanup({ failed })` | Quiesce updates and clean eligible temporary previews. Failed/partial finals and retained or promoted previews stay protected. Cleanup failure cannot authorize resending accepted content. |
| `retainPreview()` | Transfer the artifact out of automatic cleanup, for example after an accepted continuation handoff. This does not claim final delivery. |
| `reset()` | Start the next explicitly admitted turn or block generation. The owner fences stale awaited completions from the new generation. The transport still owns its corresponding message-identity rotation. |
| `reset()` | Start the next admitted turn, assistant answer, or block generation. The owner fences stale awaited completions from the new generation. The transport still owns its corresponding message-identity rotation; advance both at the assistant boundary. |
The read-only `finalStarted`, `finalDelivered`, `finalSucceeded`, `finalFailed`,
`finalSuppressed`, and `previewFinalized` properties are projections of that owner.

View file

@ -19,6 +19,7 @@ import {
ingestDraftLaneSegments,
prepareQueuedAnswerBlock,
repositionLaneForNewMessage,
resetLaneState,
rotateLaneForNewMessage,
waitForDraftEvents,
} from "./bot-message-dispatch-draft.js";
@ -250,10 +251,18 @@ export async function runTelegramDispatchTurn(turn: Turn) {
? () => {
const queued = enqueueDraftEvent(turn, async () => {
resetReasoningStepState(turn);
const previousAnswerDelivered = turn.previewLifecycle.finalDelivered;
turn.previewLifecycle.reset();
turn.finalDispatchClaimed = false;
turn.progressCompositor.beginAssistantMessage();
if (turn.answerLane.finalized) {
await rotateLaneForNewMessage(turn, turn.answerLane);
turn.rotateAnswerLaneWhenQueuedBlocksSettle = false;
} else if (previousAnswerDelivered) {
// A fresh final may have used the durable sender without leaving a draft ID.
turn.answerLane.stream?.forceNewMessage();
resetLaneState(turn, turn.answerLane);
turn.rotateAnswerLaneWhenQueuedBlocksSettle = false;
} else if (
turn.answerLane.hasStreamedMessage &&
!turn.activeAnswerDraftIsToolProgressOnly &&

View file

@ -52,23 +52,32 @@ async function setupObservedProgressTransport() {
text,
};
});
vi.spyOn(bot.api, "editMessageText").mockImplementation(async (_chatId, messageId, text) => {
if (typeof text !== "string") {
throw new Error("This transport fixture expects legacy Telegram text.");
}
if (!messages.has(messageId)) {
throw new Error("Bad Request: message to edit not found");
}
messages.set(messageId, text);
return true;
});
const editMessageText = vi
.spyOn(bot.api, "editMessageText")
.mockImplementation(async (_chatId, messageId, text) => {
if (typeof text !== "string") {
throw new Error("This transport fixture expects legacy Telegram text.");
}
if (!messages.has(messageId)) {
throw new Error("Bad Request: message to edit not found");
}
messages.set(messageId, text);
return true;
});
const deleteMessage = vi
.spyOn(bot.api, "deleteMessage")
.mockImplementation(async (_chatId, messageId) => {
messages.delete(messageId);
return true;
});
return { bot, messages, sendMessage, deleteMessage, deliver: delivery.deliverStructuredReplies };
return {
bot,
messages,
sendMessage,
editMessageText,
deleteMessage,
deliver: delivery.deliverStructuredReplies,
};
}
function progressContext(
@ -422,4 +431,64 @@ describeTelegramDispatch("dispatchTelegramMessage final-delivery-lifecycle", ()
vi.useRealTimers();
}
});
it.each(["accepted", "rejected"] as const)(
"uses the second assistant preview as the current final when delivery is %s",
async (outcome) => {
vi.useFakeTimers();
try {
const { bot, messages, sendMessage, editMessageText } =
await setupObservedProgressTransport();
const status = createStatusReactionController();
const first = "First accepted answer";
const partial = "Second answer in progress while inspecting the result";
const second = "Second complete answer";
dispatchReplyWithBufferedBlockDispatcher.mockImplementationOnce(async (params) => {
await params.dispatcherOptions.deliver({ text: first }, { kind: "final" });
await params.replyOptions?.onAssistantMessageStart?.();
await params.replyOptions?.onPartialReply?.({ text: partial });
await createTelegramDraftStream.mock.results[0]?.value?.flush();
expect([...messages.values()]).toEqual([first, partial]);
if (outcome === "rejected") {
editMessageText.mockRejectedValueOnce(new Error("second final edit rejected"));
sendMessage.mockRejectedValueOnce(new Error("second final send rejected"));
}
return dispatchThroughSharedOwner({
...params,
replyResolver: async () => ({ text: second }),
});
});
await dispatchWithContext({
bot,
cfg: { channels: { telegram: { botToken: "synthetic-test-token" } } },
context: progressContext(status),
streamMode: "partial",
telegramCfg: { streaming: { mode: "partial" } },
retryDispatchErrors: true,
suppressFailureFallback: true,
});
await vi.advanceTimersByTimeAsync(5_000);
const visible = [...messages.entries()];
expect(visible[0]?.[1]).toBe(first);
expect(sendMessage.mock.calls.filter(([, text]) => text === first)).toHaveLength(1);
if (outcome === "accepted") {
expect(visible).toHaveLength(2);
expect(visible[1]?.[1]).toBe(second);
expect(sendMessage.mock.calls.filter(([, text]) => text === second)).toHaveLength(0);
expect(status.setDone).toHaveBeenCalledOnce();
expect(status.setError).not.toHaveBeenCalled();
} else {
expect(visible.some(([, text]) => text === second)).toBe(false);
expect(visible.some(([, text]) => text.includes("OpenClaw chat history"))).toBe(true);
expect(sendMessage.mock.calls.filter(([, text]) => text === second)).toHaveLength(1);
expect(status.setError).toHaveBeenCalledOnce();
expect(status.setDone).not.toHaveBeenCalled();
}
} finally {
vi.clearAllTimers();
vi.useRealTimers();
}
},
);
});