diff --git a/packages/integrations/src/google.ts b/packages/integrations/src/google.ts index c94e84256..29eeb66b3 100644 --- a/packages/integrations/src/google.ts +++ b/packages/integrations/src/google.ts @@ -15,6 +15,14 @@ const GMAIL = "https://gmail.googleapis.com/gmail/v1/users/me"; const CALENDAR = "https://www.googleapis.com/calendar/v3"; export const MAX_ATTACHMENT_BYTES = 10 * 1024 * 1024; export const MAX_TOTAL_ATTACHMENT_BYTES = 20 * 1024 * 1024; +export const MAX_READ_RETRIES = 2; +export const DEFAULT_RETRY_DELAY_MS = 1000; +export const RETRYABLE_READ_STATUS_CODES = new Set([429, 500, 502, 503, 504]); +const RETRYABLE_READ_403_REASONS = new Set(["rateLimitExceeded", "userRateLimitExceeded"]); + +export function isRetryableReadStatus(status: number): boolean { + return RETRYABLE_READ_STATUS_CODES.has(status); +} const MAX_JSON_BYTES = Math.ceil((MAX_ATTACHMENT_BYTES * 4) / 3) + 1024 * 1024; export class OutcomeUnknownError extends Error { @@ -433,6 +441,32 @@ function validateAttachments(attachments: MailAttachment[]): void { throw new Error("Total attachment size exceeds the 20 MiB limit"); } +function parseRetryAfter(header: string | null): number | undefined { + if (!header) return undefined; + const seconds = Number(header); + if (Number.isFinite(seconds) && seconds >= 0) { + return Math.min(Math.round(seconds * 1000), 30_000); + } + const date = Date.parse(header); + if (Number.isFinite(date)) { + const delay = date - Date.now(); + return delay > 0 ? Math.min(delay, 30_000) : 0; + } + return undefined; +} + +function retryDelay(attempt: number): number { + const base = DEFAULT_RETRY_DELAY_MS * 2 ** attempt; + return base + Math.floor(Math.random() * base); +} + +const googleErrorSchema = z.object({ + error: z.object({ + message: z.string().optional(), + errors: z.array(z.object({ reason: z.string().optional() })).optional(), + }), +}); + async function readJson(response: Response): Promise { if (Number(response.headers.get("content-length")) > MAX_JSON_BYTES) { await response.body?.cancel(); @@ -462,9 +496,15 @@ async function readJson(response: Response): Promise { export class GoogleClient { private readonly fetcher: typeof fetch; private readonly getAccessToken: () => Promise; - constructor(options: { getAccessToken: () => Promise; fetch?: typeof fetch }) { + private readonly sleep: (ms: number) => Promise; + constructor(options: { + getAccessToken: () => Promise; + fetch?: typeof fetch; + sleep?: (ms: number) => Promise; + }) { this.fetcher = options.fetch ?? fetch; this.getAccessToken = options.getAccessToken; + this.sleep = options.sleep ?? ((ms) => new Promise((resolve) => setTimeout(resolve, ms))); } private async mapMessage(message: z.infer): Promise { @@ -593,50 +633,73 @@ export class GoogleClient { const token = await this.getAccessToken(); if (!token || /[\r\n]/.test(token)) throw new Error("Google access token is missing or invalid; reconnect Google"); - let response: Response; - try { - response = await this.fetcher(url, { - method, - headers: { - Authorization: `Bearer ${token}`, - Accept: "application/json", - ...(body === undefined ? {} : { "Content-Type": "application/json" }), - ...conditionalHeaders, - }, - body: body === undefined ? undefined : JSON.stringify(body), - signal: AbortSignal.timeout(30000), - redirect: "error", - }); - } catch { - if (write) throw new OutcomeUnknownError(); - throw new Error("Could not reach Google; check the connection and try again"); - } - if (write && (response.status >= 500 || response.status === 408)) { + + const maxAttempts = write ? 1 : 1 + MAX_READ_RETRIES; + for (let attempt = 0; attempt < maxAttempts; attempt++) { + let response: Response; try { - await response.body?.cancel(); + response = await this.fetcher(url, { + method, + headers: { + Authorization: `Bearer ${token}`, + Accept: "application/json", + ...(body === undefined ? {} : { "Content-Type": "application/json" }), + ...conditionalHeaders, + }, + body: body === undefined ? undefined : JSON.stringify(body), + signal: AbortSignal.timeout(30000), + redirect: "error", + }); } catch { + if (write) throw new OutcomeUnknownError(); + if (attempt + 1 < maxAttempts) { + await this.sleep(retryDelay(attempt)); + continue; + } + throw new Error("Could not reach Google; check the connection and try again"); + } + if (write && (response.status >= 500 || response.status === 408)) { + try { + await response.body?.cancel(); + } catch { + throw new OutcomeUnknownError(); + } throw new OutcomeUnknownError(); } - throw new OutcomeUnknownError(); - } - if (!response.ok) { - let detail = response.statusText || "Request failed"; + if (!response.ok) { + let detail = response.statusText || "Request failed"; + let retryableReason = false; + try { + const result = googleErrorSchema.safeParse(await readJson(response)); + if (result.success) { + detail = (result.data.error.message ?? detail).slice(0, 500); + retryableReason = + response.status === 403 && + (result.data.error.errors ?? []).some((error) => + RETRYABLE_READ_403_REASONS.has(error.reason ?? ""), + ); + } + } catch { + throw new GoogleApiError(response.status, detail); + } + if ( + !write && + (isRetryableReadStatus(response.status) || retryableReason) && + attempt + 1 < maxAttempts + ) { + const retryAfter = parseRetryAfter(response.headers.get("retry-after")); + await this.sleep(retryAfter ?? retryDelay(attempt)); + continue; + } + throw new GoogleApiError(response.status, detail); + } + if (method === "DELETE" && response.status === 204) return undefined; try { - const result = z - .object({ error: z.object({ message: z.string() }) }) - .safeParse(await readJson(response)); - if (result.success) detail = result.data.error.message.slice(0, 500); + return await readJson(response); } catch { - /* Preserve the definite HTTP rejection even if its body is not JSON. */ + if (write) throw new OutcomeUnknownError(); + throw new Error("Google returned an invalid or oversized response"); } - throw new GoogleApiError(response.status, detail); - } - if (method === "DELETE" && response.status === 204) return undefined; - try { - return await readJson(response); - } catch { - if (write) throw new OutcomeUnknownError(); - throw new Error("Google returned an invalid or oversized response"); } } diff --git a/tests/google.test.ts b/tests/google.test.ts index c5e607b23..14e59488c 100644 --- a/tests/google.test.ts +++ b/tests/google.test.ts @@ -2,20 +2,32 @@ import assert from "node:assert/strict"; import test from "node:test"; import { emailDraftSchema, eventDraftSchema } from "../packages/domain/src/index.ts"; import { + DEFAULT_RETRY_DELAY_MS, GoogleApiError, GoogleClient, + isRetryableReadStatus, MAX_ATTACHMENT_BYTES, + MAX_READ_RETRIES, MAX_TOTAL_ATTACHMENT_BYTES, OutcomeUnknownError, + RETRYABLE_READ_STATUS_CODES, } from "../packages/integrations/src/google.ts"; -function clientWith(handler: (request: Request) => Response | Promise) { +function clientWith( + handler: (request: Request) => Response | Promise, + sleep?: (ms: number) => Promise, +) { return new GoogleClient({ getAccessToken: async () => "synthetic-access-token", fetch: async (input, init) => handler(new Request(input, init)), + sleep, }); } const json = (data: unknown, status = 200) => Response.json(data, { status }); +function assertJitteredDelay(actual: number, base: number) { + assert.ok(actual >= base, `${actual} should be at least ${base}`); + assert.ok(actual < base * 2, `${actual} should be below ${base * 2}`); +} const email = () => emailDraftSchema.parse({ to: ["reader@example.com"], @@ -920,6 +932,242 @@ test("single-event validation is repeated at execution and read failures never d assert.equal(writes, 0); }); +test("GET reads retry on HTTP 429 rate limit with Retry-After header and succeed", async () => { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + if (attempts === 1) { + return new Response(JSON.stringify({ error: { message: "Rate limit exceeded" } }), { + status: 429, + headers: { "Content-Type": "application/json", "Retry-After": "3" }, + }); + } + return json({ + items: [{ id: "c1", summary: "Personal", timeZone: "UTC", accessRole: "owner" }], + }); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + const calendars = await client.listCalendars(); + assert.equal(attempts, 2); + assert.deepEqual(sleeps, [3000]); + assert.equal(calendars.length, 1); +}); + +test("GET reads retry Google 403 rate-limit reasons with jittered backoff", async () => { + for (const reason of ["rateLimitExceeded", "userRateLimitExceeded"]) { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + if (attempts === 1) { + return json({ error: { message: "Slow down", errors: [{ reason }] } }, 403); + } + return json({ + items: [{ id: "c1", summary: "Personal", timeZone: "UTC", accessRole: "owner" }], + }); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + const calendars = await client.listCalendars(); + assert.equal(attempts, 2); + assert.equal(sleeps.length, 1); + assertJitteredDelay(sleeps[0], DEFAULT_RETRY_DELAY_MS); + assert.equal(calendars.length, 1); + } +}); + +test("GET reads do not retry when the error response body is malformed", async () => { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + return new Response("{not json", { + status: 503, + headers: { "Content-Type": "application/json" }, + }); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + await assert.rejects( + client.listCalendars(), + (error: unknown) => error instanceof GoogleApiError && error.status === 503, + ); + assert.equal(attempts, 1); + assert.deepEqual(sleeps, []); +}); + +test("GET reads retry transient 503 and back off up to MAX_READ_RETRIES before failing", async () => { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + return new Response( + JSON.stringify({ error: { message: "Backend temporarily unavailable" } }), + { + status: 503, + headers: { "Content-Type": "application/json" }, + }, + ); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + await assert.rejects( + client.listCalendars(), + (error: unknown) => + error instanceof GoogleApiError && + error.status === 503 && + /Backend temporarily unavailable/.test(error.message), + ); + assert.equal(attempts, 1 + MAX_READ_RETRIES); + assert.equal(sleeps.length, 2); + assertJitteredDelay(sleeps[0], DEFAULT_RETRY_DELAY_MS); + assertJitteredDelay(sleeps[1], DEFAULT_RETRY_DELAY_MS * 2); +}); + +test("GET reads retry transient network error and recover", async () => { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + if (attempts === 1) throw new TypeError("network socket disconnected"); + return json({ items: [eventResponse] }); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + const events = await client.listEvents(); + assert.equal(attempts, 2); + assert.equal(sleeps.length, 1); + assertJitteredDelay(sleeps[0], DEFAULT_RETRY_DELAY_MS); + assert.equal(events.length, 1); +}); + +test("non-retryable client errors on GET fail immediately without retrying", async () => { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + return new Response(JSON.stringify({ error: { message: "Calendar not found" } }), { + status: 404, + headers: { "Content-Type": "application/json" }, + }); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + await assert.rejects( + client.listCalendars(), + (error: unknown) => error instanceof GoogleApiError && error.status === 404, + ); + assert.equal(attempts, 1); + assert.equal(sleeps.length, 0); +}); + +test("isRetryableReadStatus predicate explicitly covers transient server errors and rate limits", () => { + for (const status of [429, 500, 502, 503, 504]) { + assert.equal(isRetryableReadStatus(status), true); + } + for (const status of [400, 401, 403, 404, 408, 422, 501]) { + assert.equal(isRetryableReadStatus(status), false); + } + assert.deepEqual(Array.from(RETRYABLE_READ_STATUS_CODES).sort(), [429, 500, 502, 503, 504]); +}); + +test("GET reads retry transient 502 Bad Gateway and succeed on retry", async () => { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + if (attempts === 1) { + return new Response(JSON.stringify({ error: { message: "Bad Gateway" } }), { + status: 502, + headers: { "Content-Type": "application/json" }, + }); + } + return json({ + items: [{ id: "c1", summary: "Personal", timeZone: "UTC", accessRole: "owner" }], + }); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + const calendars = await client.listCalendars(); + assert.equal(attempts, 2); + assert.equal(sleeps.length, 1); + assertJitteredDelay(sleeps[0], DEFAULT_RETRY_DELAY_MS); + assert.equal(calendars.length, 1); +}); + +test("GET reads retry transient 504 Gateway Timeout and back off up to MAX_READ_RETRIES before failing", async () => { + let attempts = 0; + const sleeps: number[] = []; + const client = clientWith( + () => { + attempts++; + return new Response(JSON.stringify({ error: { message: "Gateway Timeout" } }), { + status: 504, + headers: { "Content-Type": "application/json" }, + }); + }, + async (ms) => { + sleeps.push(ms); + }, + ); + await assert.rejects( + client.listCalendars(), + (error: unknown) => + error instanceof GoogleApiError && + error.status === 504 && + /Gateway Timeout/.test(error.message), + ); + assert.equal(attempts, 1 + MAX_READ_RETRIES); + assert.equal(sleeps.length, 2); + assertJitteredDelay(sleeps[0], DEFAULT_RETRY_DELAY_MS); + assertJitteredDelay(sleeps[1], DEFAULT_RETRY_DELAY_MS * 2); +}); + +test("write operations on all transient 5xx codes remain strictly single-attempt without retrying", async () => { + for (const status of [500, 502, 503, 504]) { + let writes = 0; + const client = clientWith((request) => { + if (request.url.endsWith("/profile")) return json({ emailAddress: "me@example.com" }); + if (request.method === "GET") return json(eventResponse); + writes++; + return new Response( + JSON.stringify({ error: { message: `transient write error ${status}` } }), + { + status, + headers: { "Content-Type": "application/json" }, + }, + ); + }); + await assert.rejects(client.sendEmail(email(), []), OutcomeUnknownError); + await assert.rejects(client.createEvent(event()), OutcomeUnknownError); + await assert.rejects(client.updateEvent("event-1", event()), OutcomeUnknownError); + assert.equal(writes, 3); + } +}); + test("mail parsing tolerates unknown charsets, RFC 2231 words, and display names with addresses", async () => { const client = clientWith((request) => new URL(request.url).pathname.endsWith("/messages")