Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 51 additions & 21 deletions src/core/longpoll.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,53 @@ function retryAfterMs(err: unknown): number | undefined {
return seconds === undefined ? undefined : seconds * 1000;
}

/** How the loop should react to a getUpdates error: `wait` ms before re-polling, plus the updated conflict count. */
type RetryPlan = { wait: number; conflicts: number };

/** Arguments to `planRetry` / `recover`: the error, the current conflict streak, the resolved retry options, and the signal. */
type RetryContext = {
err: unknown;
conflicts: number;
retry: boolean;
retryDelayMs: number;
conflictRetryDelayMs: number;
maxConflictRetries: number;
onError?: (err: unknown) => void;
signal?: AbortSignal;
};

/** Decide how to handle a getUpdates failure. Throws the error to stop the loop
* (non-retryable, or the conflict budget is exhausted); otherwise returns the
* wait before re-polling and the new consecutive-conflict count. */
function planRetry({ err, conflicts, retry, retryDelayMs, conflictRetryDelayMs, maxConflictRetries, onError }: RetryContext): RetryPlan {
const pollConflict = isPollConflict(err);
if (!retry || !(isTransientError(err) || pollConflict)) throw err;
// A conflict advances its own bounded counter; any other transient breaks the streak.
const next = pollConflict ? conflicts + 1 : 0;
if (pollConflict && next > maxConflictRetries) throw err;
onError?.(err);
// A conflict waits its own longer delay; otherwise honor `retry_after`
// (e.g. a 429 flood-wait) when present, else the default delay.
const wait = pollConflict ? conflictRetryDelayMs : (retryAfterMs(err) ?? retryDelayMs);
log("getUpdates %s; retry in %dms", pollConflict ? `conflict ${next}/${maxConflictRetries}` : "failed", Math.round(wait));
return { wait, conflicts: next };
}

/** Recover from a getUpdates failure: wait `plan.wait` ms then resume with the new
* conflict count, or `"stop"` the loop (the signal aborted, before or during the
* wait). `planRetry` may throw here to stop the loop on a non-retryable error. */
async function recover(ctx: RetryContext): Promise<{ conflicts: number } | "stop"> {
const { signal } = ctx;
if (signal?.aborted) return "stop"; // cancelled - swallow the abort error
const plan = planRetry(ctx);
try {
await delay(plan.wait, signal);
} catch {
return "stop"; // aborted during the wait
}
return { conflicts: plan.conflicts };
}

/** Async-generator update source (ADR-004): long-polls `getUpdates` and yields each update until the signal aborts. */
export async function* longPoll(api: Api, options: LongPollOptions = {}, signal?: AbortSignal): AsyncGenerator<Update> {
let offset = options.offset;
Expand Down Expand Up @@ -62,29 +109,12 @@ export async function* longPoll(api: Api, options: LongPollOptions = {}, signal?
signal,
);
} catch (err) {
// cancelled - swallow the abort error
if (signal?.aborted) return;
// A 409 (another instance polling the same token) is transient for polling
// but bounded, so an overlapping redeploy heals while a real two-instance
// deployment still surfaces. Everything else uses `isTransientError`.
const pollConflict = isPollConflict(err);
if (!retry || !(isTransientError(err) || pollConflict)) throw err;
if (pollConflict) {
if (++conflicts > maxConflictRetries) throw err;
} else {
conflicts = 0; // a non-conflict transient breaks the *consecutive*-conflict streak
}
onError?.(err);
// A conflict waits its own longer delay; otherwise honor `retry_after`
// (e.g. a 429 flood-wait) when present, else the default delay.
const wait = pollConflict ? conflictRetryDelayMs : (retryAfterMs(err) ?? retryDelayMs);
log("getUpdates %s; retry in %dms", pollConflict ? `conflict ${conflicts}/${maxConflictRetries}` : "failed", Math.round(wait));
try {
await delay(wait, signal);
} catch {
// aborted during the wait
return;
}
// deployment still surfaces. `recover` throws for anything non-retryable.
const outcome = await recover({ err, conflicts, retry, retryDelayMs, conflictRetryDelayMs, maxConflictRetries, onError, signal });
if (outcome === "stop") return;
conflicts = outcome.conflicts;
// retry WITHOUT advancing offset
continue;
}
Expand Down