Skip to content

Commit dc09324

Browse files
committed
fix(telegram): narrow live claim recovery
1 parent 3be0fe7 commit dc09324

2 files changed

Lines changed: 19 additions & 14 deletions

File tree

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1976,7 +1976,7 @@ describe("TelegramPollingSession", () => {
19761976
});
19771977
});
19781978

1979-
it("fails timed-out live-owned claims before draining later same-lane updates", async () => {
1979+
it("fails timed-out current-process claims before draining later same-lane updates", async () => {
19801980
await withTempSpool(async (tempDir) => {
19811981
const abort = new AbortController();
19821982
const log = vi.fn();
@@ -2022,7 +2022,7 @@ describe("TelegramPollingSession", () => {
20222022
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]);
20232023
expectLogIncludes(
20242024
log,
2025-
"spooled update 42 Telegram spooled update claim held by a live worker",
2025+
"spooled update 42 Telegram spooled update claim owned by this process",
20262026
);
20272027
stopWorker();
20282028
});

extensions/telegram/src/polling-session.ts

Lines changed: 17 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -658,31 +658,36 @@ export class TelegramPollingSession {
658658
return deferredSpooledUpdateClaimsByKey.has(buildDeferredSpooledUpdateClaimKey(update));
659659
}
660660

661-
#isTimedOutSpooledUpdateClaim(update: ClaimedTelegramSpooledUpdate): boolean {
662-
const claimedAt = update.claim?.claimedAt;
663-
return claimedAt !== undefined && Date.now() - claimedAt >= this.#spooledUpdateHandlerTimeoutMs;
664-
}
665-
666-
async #failTimedOutLiveOwnedSpooledUpdateClaims(params: {
661+
async #failTimedOutCurrentProcessSpooledUpdateClaims(params: {
667662
activeLaneKeys: Set<string>;
668663
spoolDir: string;
669664
}): Promise<void> {
670665
const claims = await listTelegramSpooledUpdateClaims({ spoolDir: params.spoolDir });
666+
const now = Date.now();
671667
for (const claim of claims) {
668+
const claimOwner = claim.claim;
669+
if (!claimOwner) {
670+
continue;
671+
}
672672
if (this.#isDeferredSpooledUpdateClaim(claim)) {
673673
continue;
674674
}
675675
if (params.activeLaneKeys.has(this.#spooledUpdateLaneKey(claim))) {
676676
continue;
677677
}
678-
if (!this.#isTimedOutSpooledUpdateClaim(claim)) {
678+
if (now - claimOwner.claimedAt < this.#spooledUpdateHandlerTimeoutMs) {
679679
continue;
680680
}
681681
if (!isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess(claim)) {
682682
continue;
683683
}
684-
const claimedForMs = Date.now() - (claim.claim?.claimedAt ?? Date.now());
685-
const message = `Telegram spooled update claim held by a live worker for ${formatDurationPrecise(claimedForMs)} without active handler state; marking failed so the lane can continue.`;
684+
// Same PID with a stale owner id means this process orphaned a previous
685+
// local handler state; different live PIDs may still be processing.
686+
if (claimOwner.processPid !== process.pid) {
687+
continue;
688+
}
689+
const claimedForMs = now - claimOwner.claimedAt;
690+
const message = `Telegram spooled update claim owned by this process for ${formatDurationPrecise(claimedForMs)} without active handler state; marking failed so the lane can continue.`;
686691
try {
687692
const failed = await failTelegramSpooledUpdateClaim({
688693
update: claim,
@@ -691,13 +696,13 @@ export class TelegramPollingSession {
691696
});
692697
if (!failed) {
693698
this.opts.log(
694-
`[telegram][diag] spooled update ${claim.updateId} live-owned claim no longer had a processing marker to fail.`,
699+
`[telegram][diag] spooled update ${claim.updateId} current-process claim no longer had a processing marker to fail.`,
695700
);
696701
continue;
697702
}
698703
} catch (err) {
699704
this.opts.log(
700-
`[telegram][diag] spooled update ${claim.updateId} live-owned claim could not be marked failed: ${formatErrorMessage(err)}`,
705+
`[telegram][diag] spooled update ${claim.updateId} current-process claim could not be marked failed: ${formatErrorMessage(err)}`,
701706
);
702707
continue;
703708
}
@@ -828,7 +833,7 @@ export class TelegramPollingSession {
828833
spoolDir: string;
829834
}): Promise<SpooledUpdateDrainResult> {
830835
const activeLaneKeys = this.#activeSpooledUpdateLaneKeysForSpool(params.spoolDir);
831-
await this.#failTimedOutLiveOwnedSpooledUpdateClaims({
836+
await this.#failTimedOutCurrentProcessSpooledUpdateClaims({
832837
activeLaneKeys,
833838
spoolDir: params.spoolDir,
834839
});

0 commit comments

Comments
 (0)