Skip to content

Commit d03ceb9

Browse files
committed
fix telegram rollback failure retry
1 parent 7a283a0 commit d03ceb9

6 files changed

Lines changed: 138 additions & 11 deletions

File tree

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,2 +1,2 @@
1-
fbb605dd3077cb40826330ee3120e1955824b9f4b10acb908e654e333ed1a1ff plugin-sdk-api-baseline.json
2-
407f0fd91ea56310b75b332f533ac3e76e769b0ea02d8e1977fa550c348ac26c plugin-sdk-api-baseline.jsonl
1+
ff7bd86cb1b243e0c94fdf9a74e7f985a7d73685b2b0cd0a8761972d145ca7a5 plugin-sdk-api-baseline.json
2+
a65283a99e28a300adffa26ed171a3e8b215d9c95e8a1656fc5ae8fd7fc011c6 plugin-sdk-api-baseline.jsonl

extensions/telegram/src/bot-handlers.runtime.ts

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1221,8 +1221,13 @@ export const registerTelegramHandlers = ({
12211221
dispatchDedupeKeys?: string[];
12221222
}): Promise<TelegramMessageProcessingResult> => {
12231223
let dispatchDedupeCommitted = false;
1224+
let dispatchDedupeRollbackAttempted = false;
12241225
const spooledReplay =
12251226
params.options?.spooledReplay === true || isTelegramSpooledReplayUpdate(params.ctx.update);
1227+
const forgetCommittedDispatchDedupeKeys = async () => {
1228+
dispatchDedupeRollbackAttempted = true;
1229+
await forgetDispatchDedupeKeys(params.dispatchDedupeKeys ?? []);
1230+
};
12261231
try {
12271232
const replyChainNodes = await buildReplyChainForMessage(params.msg);
12281233
const { replyMedia, replyChain } = await resolveReplyMediaForChain(
@@ -1252,14 +1257,14 @@ export const registerTelegramHandlers = ({
12521257
if (result.kind === "completed" && !dispatchDedupeCommitted) {
12531258
await commitDispatchDedupeKeys(params.dispatchDedupeKeys ?? []);
12541259
} else if (result.kind === "failed-retryable" && dispatchDedupeCommitted && spooledReplay) {
1255-
await forgetDispatchDedupeKeys(params.dispatchDedupeKeys ?? []);
1260+
await forgetCommittedDispatchDedupeKeys();
12561261
} else if (result.kind !== "completed" && !dispatchDedupeCommitted) {
12571262
releaseDispatchDedupeKeys(params.dispatchDedupeKeys ?? []);
12581263
}
12591264
return result;
12601265
} catch (err) {
1261-
if (dispatchDedupeCommitted && spooledReplay) {
1262-
await forgetDispatchDedupeKeys(params.dispatchDedupeKeys ?? []);
1266+
if (dispatchDedupeCommitted && spooledReplay && !dispatchDedupeRollbackAttempted) {
1267+
await forgetCommittedDispatchDedupeKeys();
12631268
} else if (!dispatchDedupeCommitted) {
12641269
releaseDispatchDedupeKeys(params.dispatchDedupeKeys ?? [], err);
12651270
}

extensions/telegram/src/message-dispatch-dedupe.test.ts

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,8 +11,11 @@ import {
1111
claimTelegramMessageDispatchReplay,
1212
commitTelegramMessageDispatchReplay,
1313
createTelegramMessageDispatchReplayGuard,
14+
forgetTelegramMessageDispatchReplay,
1415
releaseTelegramMessageDispatchReplay,
1516
TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE,
17+
TelegramMessageDispatchReplayForgetError,
18+
type TelegramMessageDispatchReplayGuard,
1619
} from "./message-dispatch-dedupe.js";
1720

1821
const tempDirs: string[] = [];
@@ -204,4 +207,27 @@ describe("Telegram message dispatch replay guard", () => {
204207
key: first.key,
205208
});
206209
});
210+
211+
it("fails rollback when a committed dispatch key cannot be forgotten", async () => {
212+
const guard = {
213+
claim: async () => ({ kind: "claimed" }),
214+
commit: async () => true,
215+
forget: async (key: string) => key !== "failed-key",
216+
hasRecent: async () => false,
217+
warmup: async () => 0,
218+
clearMemory: () => {},
219+
memorySize: () => 0,
220+
release: () => {},
221+
} satisfies TelegramMessageDispatchReplayGuard;
222+
223+
await expect(
224+
forgetTelegramMessageDispatchReplay({
225+
guard,
226+
keys: ["ok-key", "failed-key", "failed-key"],
227+
}),
228+
).rejects.toMatchObject({
229+
name: TelegramMessageDispatchReplayForgetError.name,
230+
failures: [{ key: "failed-key" }],
231+
});
232+
});
207233
});

extensions/telegram/src/message-dispatch-dedupe.ts

Lines changed: 43 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,32 @@ export type TelegramMessageDispatchClaim =
1919
| { kind: "duplicate" }
2020
| { kind: "invalid" };
2121

22+
export type TelegramMessageDispatchReplayForgetFailure = {
23+
key: string;
24+
error?: unknown;
25+
};
26+
27+
export class TelegramMessageDispatchReplayForgetError extends Error {
28+
readonly failures: TelegramMessageDispatchReplayForgetFailure[];
29+
override readonly cause: unknown;
30+
31+
constructor(failures: readonly TelegramMessageDispatchReplayForgetFailure[]) {
32+
const count = failures.length;
33+
super(`telegram message dispatch dedupe rollback failed for ${count} key(s)`, {
34+
cause: failures.find((failure) => failure.error !== undefined)?.error,
35+
});
36+
this.name = "TelegramMessageDispatchReplayForgetError";
37+
this.failures = [...failures];
38+
this.cause = failures.find((failure) => failure.error !== undefined)?.error;
39+
}
40+
}
41+
42+
export function isTelegramMessageDispatchReplayForgetError(
43+
error: unknown,
44+
): error is TelegramMessageDispatchReplayForgetError {
45+
return error instanceof TelegramMessageDispatchReplayForgetError;
46+
}
47+
2248
function sanitizeFileSegment(value: string): string {
2349
const trimmed = value.trim();
2450
if (!trimmed) {
@@ -137,11 +163,23 @@ export async function forgetTelegramMessageDispatchReplay(params: {
137163
keys?: readonly string[];
138164
}): Promise<void> {
139165
const keys = normalizeReplayKeys(params.keys);
140-
await Promise.all(
141-
keys.map((key) =>
142-
params.guard.forget(key, { namespace: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE }),
143-
),
144-
);
166+
const failures = (
167+
await Promise.all(
168+
keys.map(async (key): Promise<TelegramMessageDispatchReplayForgetFailure | null> => {
169+
try {
170+
const forgotten = await params.guard.forget(key, {
171+
namespace: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE,
172+
});
173+
return forgotten ? null : { key };
174+
} catch (error) {
175+
return { key, error };
176+
}
177+
}),
178+
)
179+
).filter((failure): failure is TelegramMessageDispatchReplayForgetFailure => Boolean(failure));
180+
if (failures.length > 0) {
181+
throw new TelegramMessageDispatchReplayForgetError(failures);
182+
}
145183
}
146184

147185
export function releaseTelegramMessageDispatchReplay(params: {

extensions/telegram/src/polling-session.test.ts

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@ let listTelegramSpooledUpdates: typeof import("./telegram-ingress-spool.js").lis
7474
let recoverStaleTelegramSpooledUpdateClaims: typeof import("./telegram-ingress-spool.js").recoverStaleTelegramSpooledUpdateClaims;
7575
let writeTelegramSpooledUpdate: typeof import("./telegram-ingress-spool.js").writeTelegramSpooledUpdate;
7676
let createTelegramSpooledReplayDeferredParticipant: typeof import("./bot-processing-outcome.js").createTelegramSpooledReplayDeferredParticipant;
77+
let TelegramMessageDispatchReplayForgetError: typeof import("./message-dispatch-dedupe.js").TelegramMessageDispatchReplayForgetError;
7778
type TelegramMessageProcessingResult =
7879
import("./bot-processing-outcome.js").TelegramMessageProcessingResult;
7980
type TelegramSpooledReplayDeferredParticipant =
@@ -619,6 +620,7 @@ describe("TelegramPollingSession", () => {
619620
} = await import("./telegram-ingress-spool.js"));
620621
({ createTelegramSpooledReplayDeferredParticipant } =
621622
await import("./bot-processing-outcome.js"));
623+
({ TelegramMessageDispatchReplayForgetError } = await import("./message-dispatch-dedupe.js"));
622624
({
623625
beginTelegramReplyFence,
624626
buildTelegramReplyFenceLaneKey,
@@ -1444,6 +1446,56 @@ describe("TelegramPollingSession", () => {
14441446
});
14451447
});
14461448

1449+
it("dead-letters buffered spooled claims when dispatch dedupe rollback fails", async () => {
1450+
await withTempSpool(async (tempDir) => {
1451+
const abort = new AbortController();
1452+
const log = vi.fn();
1453+
const participants: TelegramSpooledReplayDeferredParticipant[] = [];
1454+
const events: string[] = [];
1455+
let attempts = 0;
1456+
await writeSpooledTestUpdates(tempDir, [topicUpdate(42, 10, "buffered rollback failure")]);
1457+
1458+
const { runPromise, stopWorker } = startIsolatedIngressSession({
1459+
abort,
1460+
spoolDir: tempDir,
1461+
log,
1462+
drainIntervalMs: 10,
1463+
handleUpdate: async (update) => {
1464+
attempts += 1;
1465+
if (attempts === 1) {
1466+
events.push(`dispatch:${update.update_id}`);
1467+
const participant = createTelegramSpooledReplayDeferredParticipant(
1468+
`test-buffer:${update.update_id}`,
1469+
);
1470+
if (!participant) {
1471+
throw new Error("expected spooled replay participant");
1472+
}
1473+
participants.push(participant);
1474+
return;
1475+
}
1476+
events.push(`duplicate-skip:${update.update_id}`);
1477+
},
1478+
});
1479+
1480+
await vi.waitFor(() => expect(participants).toHaveLength(1));
1481+
participants[0]?.settle({
1482+
kind: "failed-retryable",
1483+
error: new TelegramMessageDispatchReplayForgetError([{ key: "committed-dispatch-key" }]),
1484+
});
1485+
1486+
await vi.waitFor(async () => expect(await failedUpdateIds(tempDir)).toEqual([42]));
1487+
expect(events).toEqual(["dispatch:42"]);
1488+
expect(await pendingUpdateIds(tempDir, "all")).toEqual([]);
1489+
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]);
1490+
expectLogIncludes(log, "non-retryable dispatch-dedupe-rollback-failed");
1491+
expectLogExcludes(log, "spooled update 42 failed; keeping for retry");
1492+
1493+
abort.abort();
1494+
stopWorker();
1495+
await runPromise;
1496+
});
1497+
});
1498+
14471499
it("fails buffered spooled claims instead of requeueing when deferred processing times out", async () => {
14481500
await withTempSpool(async (tempDir) => {
14491501
const abort = new AbortController();

extensions/telegram/src/polling-session.ts

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import {
2626
} from "./bot-processing-outcome.js";
2727
import { createTelegramBot } from "./bot.js";
2828
import type { TelegramTransport } from "./fetch.js";
29+
import { isTelegramMessageDispatchReplayForgetError } from "./message-dispatch-dedupe.js";
2930
import { isRecoverableTelegramNetworkError } from "./network-errors.js";
3031
import { TelegramPollingLivenessTracker } from "./polling-liveness.js";
3132
import { createTelegramPollingStatusPublisher } from "./polling-status.js";
@@ -141,7 +142,7 @@ function normalizeTelegramAccountId(accountId?: string | null): string {
141142
}
142143

143144
type NonRetryableSpooledUpdateFailure = {
144-
reason: "missing-agent-harness";
145+
reason: "missing-agent-harness" | "dispatch-dedupe-rollback-failed";
145146
message: string;
146147
};
147148

@@ -153,6 +154,11 @@ function resolveNonRetryableSpooledUpdateFailure(
153154
current.error,
154155
])) {
155156
const message = formatErrorMessage(candidate);
157+
if (isTelegramMessageDispatchReplayForgetError(candidate)) {
158+
// A committed dispatch key that cannot be rolled back makes retry unsafe:
159+
// the next replay can be duplicate-suppressed and then deleted.
160+
return { reason: "dispatch-dedupe-rollback-failed", message };
161+
}
156162
if (
157163
readErrorName(candidate) === MISSING_AGENT_HARNESS_ERROR_NAME ||
158164
MISSING_AGENT_HARNESS_MESSAGE_RE.test(message)

0 commit comments

Comments
 (0)