Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
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
167 changes: 167 additions & 0 deletions connectors/gmail/src/gmail-api-retry.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,167 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import { GmailApi, GmailApiError, isGmailRateLimitError } from "./gmail-api";

/**
* Rate-limit / transient retry for {@link GmailApi.call}. Gmail write-backs
* (star / read via modifyThread) had NO backoff: a single per-user-per-minute
* quota 403 (`rateLimitExceeded`) threw straight through and the change was
* dropped (PostHog 019ed581). `call()` now absorbs brief blips in-process and
* throws on sustained ones so the caller can defer.
*/
describe("GmailApi.call — rate-limit / transient retry", () => {
afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
vi.restoreAllMocks();
});

it("retries a 429 then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("rate", { status: 429, statusText: "Too Many Requests" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 403 rateLimitExceeded then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response(
JSON.stringify({
error: { errors: [{ reason: "rateLimitExceeded" }], message: "Quota exceeded" },
}),
{ status: 403, statusText: "Forbidden" }
)
)
.mockResolvedValueOnce(new Response(null, { status: 204 }));
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").modifyThread("t1", ["STARRED"]);
await vi.runAllTimersAsync();
await expect(p).resolves.toBeUndefined();
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("retries a 5xx then succeeds", async () => {
vi.useFakeTimers();
const fetchMock = vi
.fn()
.mockResolvedValueOnce(
new Response("oops", { status: 503, statusText: "Service Unavailable" })
)
.mockResolvedValueOnce(
new Response(JSON.stringify({ ok: true }), {
status: 200,
headers: { "Content-Type": "application/json" },
})
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
await vi.runAllTimersAsync();
await expect(p).resolves.toEqual({ ok: true });
expect(fetchMock).toHaveBeenCalledTimes(2);
});

it("does NOT retry a 403 that is not a rate-limit (permission error)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response(JSON.stringify({ error: { message: "Insufficient Permission" } }), {
status: 403,
statusText: "Forbidden",
})
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
name: "GmailApiError",
status: 403,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry a 404", async () => {
const fetchMock = vi
.fn()
.mockResolvedValue(new Response("missing", { status: 404, statusText: "Not Found" }));
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/threads/x")).rejects.toMatchObject({
status: 404,
});
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("does NOT retry in-process when Retry-After exceeds the cap (defers instead)", async () => {
const fetchMock = vi.fn().mockResolvedValue(
new Response("rate", { status: 429, headers: { "Retry-After": "120" } })
);
vi.stubGlobal("fetch", fetchMock);

await expect(new GmailApi("tok").call("/profile")).rejects.toMatchObject({
status: 429,
});
// A 2-minute wait belongs in the deferred drain, not an in-flight isolate.
expect(fetchMock).toHaveBeenCalledTimes(1);
});

it("throws after exhausting retries on a persistent 429", async () => {
vi.useFakeTimers();
// Fresh Response per call — a Response body can only be read once.
const fetchMock = vi
.fn()
.mockImplementation(
async () =>
new Response("rate", { status: 429, statusText: "Too Many Requests" })
);
vi.stubGlobal("fetch", fetchMock);

const p = new GmailApi("tok").call("/profile");
p.catch(() => {}); // avoid unhandled-rejection noise while timers advance
await vi.runAllTimersAsync();
await expect(p).rejects.toMatchObject({ status: 429 });
expect(fetchMock).toHaveBeenCalledTimes(3); // 1 initial + 2 retries
});
});

describe("isGmailRateLimitError", () => {
it("is true for HTTP 429", () => {
expect(isGmailRateLimitError(new GmailApiError(429, "Too Many Requests", ""))).toBe(true);
});
it("is true for a 403 rateLimitExceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", '...reason: "rateLimitExceeded"...'))
).toBe(true);
});
it("is true for a 403 Quota exceeded body", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Quota exceeded for quota metric"))
).toBe(true);
});
it("is false for a 403 permission error", () => {
expect(
isGmailRateLimitError(new GmailApiError(403, "Forbidden", "Insufficient Permission"))
).toBe(false);
});
it("is false for a 404", () => {
expect(isGmailRateLimitError(new GmailApiError(404, "Not Found", ""))).toBe(false);
});
it("is false for a non-GmailApiError that merely mentions the marker", () => {
expect(isGmailRateLimitError(new Error("rateLimitExceeded"))).toBe(false);
});
});
134 changes: 119 additions & 15 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,6 +78,39 @@ export class GmailApiError extends Error {
}
}

/**
* True when a {@link GmailApiError} is a Gmail rate-limit / quota rejection:
* an HTTP 429, or a 403 carrying one of Gmail's unambiguous quota markers
* (`rateLimitExceeded` / `userRateLimitExceeded` / `Quota exceeded`). These are
* expected under load and self-resolve once the per-user-per-minute window
* clears, so callers retry/defer them rather than dropping the write-back or
* paging error tracking. Gated on the markers — NOT a bare 403 — so genuine
* permission failures (`Insufficient Permission`) still surface.
*/
export function isGmailRateLimitError(error: unknown): boolean {
if (!(error instanceof GmailApiError)) return false;
if (error.status === 429) return true;
return (
error.status === 403 &&
/rateLimitExceeded|userRateLimitExceeded|Quota exceeded/i.test(error.message)
);
}

/**
* In-process retry budget for {@link GmailApi.call}. Kept small and short so a
* brief blip (a momentary 429, a 5xx, a dropped connection) is absorbed inside
* the current execution without risking the worker's wall-clock budget. Sustained
* rate-limits exceed this and throw, so the caller (deferred write-back drain,
* incremental-sync pending list) can reschedule past the quota window.
*/
const GMAIL_CALL_MAX_ATTEMPTS = 3;
const GMAIL_CALL_BACKOFF_MS = [500, 1500];
/**
* Honor a server `Retry-After` only up to this bound; a longer wait belongs in a
* scheduled retry, not an in-flight isolate, so we throw and let the caller defer.
*/
const GMAIL_RETRY_AFTER_MAX_MS = 3000;

export class GmailApi {
private baseUrl = "https://gmail.googleapis.com/gmail/v1/users/me";

Expand DownExpand Up@@ -114,25 +147,96 @@ export class GmailApi {
"Content-Type": "application/json",
};

const response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
// Bounded in-process retry for transient failures (rate-limit / 5xx /
// dropped connection). A momentary blip is absorbed here; a sustained one
// exceeds the budget and throws so the caller can defer past the quota
// window (see deferred write-back drain / incremental-sync pending list).
let lastError: unknown;
for (let attempt = 0; attempt < GMAIL_CALL_MAX_ATTEMPTS; attempt++) {
let response: Response;
try {
response = await fetch(url.toString(), {
method,
headers,
body: body ? JSON.stringify(body) : undefined,
});
} catch (networkError) {
// fetch() rejects on a dropped/aborted connection — transient.
lastError = networkError;
if (attempt < GMAIL_CALL_MAX_ATTEMPTS - 1) {
await this.sleep(GMAIL_CALL_BACKOFF_MS[attempt] ?? 0);
continue;
}
throw networkError;
}

if (response.ok) {
// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end
// of JSON input"; this escaped through setupWatch()'s unguarded
// stopWatch() recovery path and surfaced as an unhandled twist
// exception. Read the body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
}

if (!response.ok) {
const errorText = await response.text();
throw new GmailApiError(response.status, response.statusText, errorText);
const error = new GmailApiError(
response.status,
response.statusText,
errorText
);
lastError = error;

const retryable = isGmailRateLimitError(error) || response.status >= 500;
if (!retryable || attempt >= GMAIL_CALL_MAX_ATTEMPTS - 1) {
throw error;
}

const delayMs = this.retryDelayMs(response, attempt);
// A long server-requested wait belongs in a scheduled retry, not here.
if (delayMs === null) throw error;
await this.sleep(delayMs);
}

// Some Gmail endpoints — notably users.stop (POST /stop, used by
// stopWatch) — return 204 No Content with an EMPTY body. Calling
// response.json() on an empty body throws "SyntaxError: Unexpected end of
// JSON input"; this escaped through setupWatch()'s unguarded stopWatch()
// recovery path and surfaced as an unhandled twist exception. Read the
// body as text and only parse it when non-empty.
const text = await response.text();
return text ? JSON.parse(text) : undefined;
// Unreachable in practice (the loop returns or throws), but satisfies the
// type checker and surfaces any logic error rather than returning undefined.
throw lastError instanceof Error
? lastError
: new GmailApiError(0, "Retry loop exhausted", String(lastError));
}

private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

/**
* Backoff for a retryable response. Honors a `Retry-After` header (seconds or
* HTTP date) up to {@link GMAIL_RETRY_AFTER_MAX_MS}; returns null when the
* server asks for longer than that, signalling the caller to throw and defer
* rather than block the isolate. Falls back to a fixed backoff schedule.
*/
private retryDelayMs(response: Response, attempt: number): number | null {
const retryAfter = response.headers.get("Retry-After");
if (retryAfter) {
const seconds = Number(retryAfter);
let ms: number;
if (Number.isFinite(seconds)) {
ms = seconds * 1000;
} else {
const when = Date.parse(retryAfter);
ms = Number.isNaN(when) ? NaN : when - Date.now();
}
if (Number.isFinite(ms)) {
if (ms > GMAIL_RETRY_AFTER_MAX_MS) return null;
return Math.max(0, ms);
}
}
return (
GMAIL_CALL_BACKOFF_MS[attempt] ??
GMAIL_CALL_BACKOFF_MS[GMAIL_CALL_BACKOFF_MS.length - 1]
);
}

public async getLabels(): Promise<GmailLabel[]> {
Expand Down
Loading
Loading