From 9fd29196986828c1bad92bfb918cce7f8870c389 Mon Sep 17 00:00:00 2001 From: Colin Ozanne Date: Tue, 16 Jun 2026 15:17:06 +0200 Subject: [PATCH] fix: retry turso query + add small caching (#2289) * fix: retry turso query + add small caching * fix: wrap sibling external-service reads in retryRead * refactor(services): Effect-based retry + address PR review --- apps/server/src/libs/middlewares/auth.ts | 10 +- .../src/app/(landing)/status/[id]/page.tsx | 75 +++++----- .../src/app/api/og/external-service/route.tsx | 4 +- .../web/src/lib/external-report-escalation.ts | 50 ++++--- packages/services/package.json | 1 + packages/services/src/__tests__/retry.test.ts | 85 +++++++++++ .../src/external-service-component/list.ts | 99 ++++++------- .../src/external-service-incident/list.ts | 107 +++++++------- .../src/external-service-report/read.ts | 133 ++++++++++-------- .../src/external-service/list-slugs.ts | 7 +- .../services/src/external-service/list.ts | 29 ++-- packages/services/src/index.ts | 1 + packages/services/src/retry.ts | 36 +++-- pnpm-lock.yaml | 3 + 14 files changed, 388 insertions(+), 252 deletions(-) diff --git a/apps/server/src/libs/middlewares/auth.ts b/apps/server/src/libs/middlewares/auth.ts index c4c5782a..4390e6e5 100644 --- a/apps/server/src/libs/middlewares/auth.ts +++ b/apps/server/src/libs/middlewares/auth.ts @@ -10,11 +10,7 @@ import { shouldUpdateLastUsed, verifyApiKeyHash, } from "@openstatus/db/src/utils/api-key"; -import { - isRetryableDbError, - isTransientServerError, - withBusyRetry, -} from "@openstatus/services"; +import { retryRead } from "@openstatus/services"; import { UnkeyCore } from "@unkey/api/core"; import { keysVerifyKey } from "@unkey/api/funcs/keysVerifyKey"; import type { Context, Next } from "hono"; @@ -25,10 +21,6 @@ import type { Variables } from "@/types"; const logger = getLogger("api-server"); -// Reads only: a 502 can land after a write partially applied. -const retryRead = (fn: () => Promise): Promise => - withBusyRetry(fn, (e) => isRetryableDbError(e) || isTransientServerError(e)); - export async function lookupWorkspace(workspaceId: number) { const _workspace = await retryRead(() => db.select().from(workspace).where(eq(workspace.id, workspaceId)).get(), diff --git a/apps/web/src/app/(landing)/status/[id]/page.tsx b/apps/web/src/app/(landing)/status/[id]/page.tsx index 004704c5..ea21922e 100644 --- a/apps/web/src/app/(landing)/status/[id]/page.tsx +++ b/apps/web/src/app/(landing)/status/[id]/page.tsx @@ -33,46 +33,51 @@ export async function generateMetadata(args: { params: Promise; }): Promise { const { id } = await args.params; - const service = await cachedGetExternalServiceBySlug(id); - if (!service) return { ...defaultMetadata, title: "Not Found" }; + try { + const service = await cachedGetExternalServiceBySlug(id); + if (!service) return { ...defaultMetadata, title: "Not Found" }; - const { escalated } = await getServiceEscalation(service); + const { escalated } = await getServiceEscalation(service); - const title = escalated - ? `Users reporting issues with ${service.name}. ${service.name} Status & Incidents` - : `Is ${service.name} Down? ${service.name} Status & Incidents`; - const description = escalated - ? `Users are reporting problems with ${service.name} in the last ${REPORT_WINDOW_MINUTES} minutes. Check the live ${service.name} status, uptime over the last ${HISTORY_DAYS} days, and recent ${service.name} incidents tracked by OpenStatus.` - : `Is ${service.name} down right now? Check the live ${service.name} status, uptime over the last ${HISTORY_DAYS} days, and recent ${service.name} incidents tracked by OpenStatus.`; - const canonicalUrl = `${BASE_URL}/status/${service.slug}`; - const ogImage = `${BASE_URL}/api/og/external-service?slug=${encodeURIComponent(service.slug)}`; - const indexable = service.deletedAt == null; + const title = escalated + ? `Users reporting issues with ${service.name}. ${service.name} Status & Incidents` + : `Is ${service.name} Down? ${service.name} Status & Incidents`; + const description = escalated + ? `Users are reporting problems with ${service.name} in the last ${REPORT_WINDOW_MINUTES} minutes. Check the live ${service.name} status, uptime over the last ${HISTORY_DAYS} days, and recent ${service.name} incidents tracked by OpenStatus.` + : `Is ${service.name} down right now? Check the live ${service.name} status, uptime over the last ${HISTORY_DAYS} days, and recent ${service.name} incidents tracked by OpenStatus.`; + const canonicalUrl = `${BASE_URL}/status/${service.slug}`; + const ogImage = `${BASE_URL}/api/og/external-service?slug=${encodeURIComponent(service.slug)}`; + const indexable = service.deletedAt == null; - return { - ...defaultMetadata, - title, - description, - alternates: { - canonical: canonicalUrl, - }, - robots: { - index: indexable, - follow: true, - }, - openGraph: { - ...ogMetadata, - title, - description, - url: canonicalUrl, - images: [ogImage], - }, - twitter: { - ...twitterMetadata, + return { + ...defaultMetadata, title, description, - images: [ogImage], - }, - }; + alternates: { + canonical: canonicalUrl, + }, + robots: { + index: indexable, + follow: true, + }, + openGraph: { + ...ogMetadata, + title, + description, + url: canonicalUrl, + images: [ogImage], + }, + twitter: { + ...twitterMetadata, + title, + description, + images: [ogImage], + }, + }; + } catch (err) { + console.warn(`[status metadata] fallback for ${id}:`, err); + return defaultMetadata; + } } export default async function Page(args: { params: Promise }) { diff --git a/apps/web/src/app/api/og/external-service/route.tsx b/apps/web/src/app/api/og/external-service/route.tsx index e89a7c30..a2c1333e 100644 --- a/apps/web/src/app/api/og/external-service/route.tsx +++ b/apps/web/src/app/api/og/external-service/route.tsx @@ -15,7 +15,9 @@ import { cn } from "@/lib/utils"; import { SIZE } from "../utils"; -export const runtime = "edge"; +// nodejs (not edge): this route pulls in the service reads + Effect retry, which +// push the bundle past the 2 MB edge limit. Node has no such cap. +export const runtime = "nodejs"; export const dynamic = "force-dynamic"; const FOOTER = "openstatus.dev/status"; diff --git a/apps/web/src/lib/external-report-escalation.ts b/apps/web/src/lib/external-report-escalation.ts index 3c6b0716..be7f9de6 100644 --- a/apps/web/src/lib/external-report-escalation.ts +++ b/apps/web/src/lib/external-report-escalation.ts @@ -8,11 +8,41 @@ import { getServiceReportWindows, } from "@openstatus/services/external-service-report"; import { OSTinybird, safePipeData } from "@openstatus/tinybird"; +import { unstable_cache } from "next/cache"; import { env } from "@/env"; const tb = new OSTinybird(env.TINY_BIRD_API_KEY); +const REPORTS_REVALIDATE_SECONDS = 30; +// Distinct from the "external-services" tag on service-lookup caches so that +// service mutations don't also flush the short-lived reports cache. +const REPORTS_TAG = "external-service-reports"; + +// Keyed by serviceId only — `since` is a moving Date.now() window, so it must be +// computed inside (keying on it would make every request a cache miss). +const cachedServiceReporters = unstable_cache( + async (serviceId: number): Promise => { + const since = new Date(Date.now() - REPORT_WINDOW_MS); + const rows = await getServiceReportWindows({ + serviceIds: [serviceId], + since, + }); + return rows[0]?.reporters ?? 0; + }, + ["external-service-report:service-reporters"], + { revalidate: REPORTS_REVALIDATE_SECONDS, tags: [REPORTS_TAG] }, +); + +const cachedComponentReporters = unstable_cache( + async (serviceId: number) => { + const since = new Date(Date.now() - REPORT_WINDOW_MS); + return getComponentReportWindows({ serviceId, since }); + }, + ["external-service-report:component-reporters"], + { revalidate: REPORTS_REVALIDATE_SECONDS, tags: [REPORTS_TAG] }, +); + type EscalationInput = { id: number; slug: string; @@ -28,16 +58,9 @@ export type ServiceEscalation = { lastFetchedAt: number; }; -async function serviceReporters( - serviceId: number, - since: Date, -): Promise { +async function serviceReporters(serviceId: number): Promise { try { - const rows = await getServiceReportWindows({ - serviceIds: [serviceId], - since, - }); - return rows[0]?.reporters ?? 0; + return await cachedServiceReporters(serviceId); } catch (err) { console.error("[escalation] service reporters read failed:", err); return 0; @@ -49,14 +72,13 @@ export async function getServiceEscalation( ): Promise { const aliasSlugs = Array.isArray(service.aliases) ? service.aliases : []; const slugChain = [service.slug, ...aliasSlugs]; - const since = new Date(Date.now() - REPORT_WINDOW_MS); const [latestRes, reporters] = await Promise.all([ safePipeData( tb.externalStatusLatest({ ids: slugChain }), "externalStatusLatest (escalation)", ), - serviceReporters(service.id, since), + serviceReporters(service.id), ]); const latestRows = [...latestRes.data].sort( @@ -84,13 +106,9 @@ export async function getComponentEscalation(args: { indicator: string; status: string; }): Promise<{ indicator: string; status: string; escalated: boolean }> { - const since = new Date(Date.now() - REPORT_WINDOW_MS); let reporters = 0; try { - const rows = await getComponentReportWindows({ - serviceId: args.serviceId, - since, - }); + const rows = await cachedComponentReporters(args.serviceId); reporters = rows.find((r) => r.componentId === args.componentId)?.reporters ?? 0; } catch (err) { diff --git a/packages/services/package.json b/packages/services/package.json index e80d5420..089c0b4d 100644 --- a/packages/services/package.json +++ b/packages/services/package.json @@ -136,6 +136,7 @@ "@openstatus/theme-store": "workspace:*", "@openstatus/tinybird": "workspace:*", "@openstatus/utils": "workspace:*", + "effect": "catalog:", "zod": "catalog:" }, "devDependencies": { diff --git a/packages/services/src/__tests__/retry.test.ts b/packages/services/src/__tests__/retry.test.ts index b2a812e9..9f26dcaa 100644 --- a/packages/services/src/__tests__/retry.test.ts +++ b/packages/services/src/__tests__/retry.test.ts @@ -4,9 +4,20 @@ import { NotFoundError } from "../errors"; import { isRetryableDbError, isTransientServerError, + retryRead, withBusyRetry, } from "../retry"; +const transient502 = () => { + const err = new Error("Failed query: select ..."); + (err as Error & { cause: unknown }).cause = { + code: "SERVER_ERROR", + status: 502, + message: "SERVER_ERROR: Server returned HTTP status 502", + }; + return err; +}; + describe("isRetryableDbError", () => { test("true for SQLITE_BUSY code", () => { expect(isRetryableDbError({ code: "SQLITE_BUSY" })).toBe(true); @@ -40,6 +51,26 @@ describe("isRetryableDbError", () => { expect(isRetryableDbError(null)).toBe(false); expect(isRetryableDbError("SQLITE_BUSY")).toBe(false); }); + + test("terminates on a self-referential cause cycle", () => { + const err: { code?: string; cause?: unknown } = {}; + err.cause = err; + expect(isRetryableDbError(err)).toBe(false); + }); + + test("finds a retryable code at the deepest searched level", () => { + // wrapped 9 times -> code sits at depth 9, the last level inspected. + let chain: { code?: string; cause?: unknown } = { code: "SQLITE_BUSY" }; + for (let i = 0; i < 9; i++) chain = { cause: chain }; + expect(isRetryableDbError(chain)).toBe(true); + }); + + test("stops at MAX_CAUSE_DEPTH (code one level too deep is unreachable)", () => { + // wrapped 10 times -> code sits at depth 10, never inspected. + let chain: { code?: string; cause?: unknown } = { code: "SQLITE_BUSY" }; + for (let i = 0; i < 10; i++) chain = { cause: chain }; + expect(isRetryableDbError(chain)).toBe(false); + }); }); describe("isTransientServerError", () => { @@ -143,3 +174,57 @@ describe("withBusyRetry", () => { expect(attempts).toBe(1); }); }); + +describe("retryRead", () => { + test("retries a drizzle-wrapped transient 502 then resolves", async () => { + let attempts = 0; + const result = await retryRead(async () => { + attempts++; + if (attempts < 3) throw transient502(); + return "ok"; + }); + expect(result).toBe("ok"); + expect(attempts).toBe(3); + }); + + test("retries SQLITE_BUSY", async () => { + let attempts = 0; + const result = await retryRead(async () => { + attempts++; + if (attempts < 2) throw { code: "SQLITE_BUSY" }; + return "ok"; + }); + expect(result).toBe("ok"); + expect(attempts).toBe(2); + }); + + test("does not retry a 4xx server error", async () => { + let attempts = 0; + const original = { + code: "SERVER_ERROR", + status: 404, + message: "Server returned HTTP status 404", + }; + const promise = retryRead(async () => { + attempts++; + throw original; + }); + await expect(promise).rejects.toBe(original); + expect(attempts).toBe(1); + }); + + test("gives up after 5 attempts on a sustained 5xx", async () => { + let attempts = 0; + const original = { + code: "SERVER_ERROR", + status: 503, + message: "Server returned HTTP status 503", + }; + const promise = retryRead(async () => { + attempts++; + throw original; + }); + await expect(promise).rejects.toBe(original); + expect(attempts).toBe(5); + }); +}); diff --git a/packages/services/src/external-service-component/list.ts b/packages/services/src/external-service-component/list.ts index 5f02d9e9..0f2f8d27 100644 --- a/packages/services/src/external-service-component/list.ts +++ b/packages/services/src/external-service-component/list.ts @@ -5,6 +5,7 @@ import { externalServiceComponent } from "@openstatus/db/src/schema"; import { defaultTb } from "../context"; import { getExternalServiceBySlug } from "../external-service"; import type { GlobalReadContext } from "../external-service/internal"; +import { retryRead } from "../retry"; export type ExternalComponentListItem = { id: number; @@ -80,27 +81,29 @@ export async function listExternalComponentsByServiceId(args: { const { ctx, externalServiceId } = args; const db = ctx?.db ?? defaultDb; - const rows = await db - .select({ - id: externalServiceComponent.id, - externalServiceId: externalServiceComponent.externalServiceId, - upstreamComponentId: externalServiceComponent.upstreamComponentId, - slug: externalServiceComponent.slug, - name: externalServiceComponent.name, - description: externalServiceComponent.description, - groupName: externalServiceComponent.groupName, - position: externalServiceComponent.position, - indicator: externalServiceComponent.indicator, - status: externalServiceComponent.status, - firstSeenAt: externalServiceComponent.firstSeenAt, - }) - .from(externalServiceComponent) - .where(eq(externalServiceComponent.externalServiceId, externalServiceId)) - .orderBy( - asc(externalServiceComponent.position), - asc(externalServiceComponent.name), - ) - .all(); + const rows = await retryRead(() => + db + .select({ + id: externalServiceComponent.id, + externalServiceId: externalServiceComponent.externalServiceId, + upstreamComponentId: externalServiceComponent.upstreamComponentId, + slug: externalServiceComponent.slug, + name: externalServiceComponent.name, + description: externalServiceComponent.description, + groupName: externalServiceComponent.groupName, + position: externalServiceComponent.position, + indicator: externalServiceComponent.indicator, + status: externalServiceComponent.status, + firstSeenAt: externalServiceComponent.firstSeenAt, + }) + .from(externalServiceComponent) + .where(eq(externalServiceComponent.externalServiceId, externalServiceId)) + .orderBy( + asc(externalServiceComponent.position), + asc(externalServiceComponent.name), + ) + .all(), + ); if (rows.length === 0) return []; @@ -186,33 +189,35 @@ export async function getExternalComponentBySlug(args: { const service = await getExternalServiceBySlug({ ctx, slug: serviceSlug }); if (!service) return { service: null, component: null }; - const rows = await db - .select({ - id: externalServiceComponent.id, - externalServiceId: externalServiceComponent.externalServiceId, - upstreamComponentId: externalServiceComponent.upstreamComponentId, - slug: externalServiceComponent.slug, - aliases: externalServiceComponent.aliases, - name: externalServiceComponent.name, - description: externalServiceComponent.description, - groupName: externalServiceComponent.groupName, - position: externalServiceComponent.position, - indicator: externalServiceComponent.indicator, - status: externalServiceComponent.status, - firstSeenAt: externalServiceComponent.firstSeenAt, - }) - .from(externalServiceComponent) - .where( - and( - eq(externalServiceComponent.externalServiceId, service.id), - or( - eq(externalServiceComponent.slug, componentSlug), - sql`EXISTS (SELECT 1 FROM json_each(${externalServiceComponent.aliases}) WHERE value = ${componentSlug})`, + const rows = await retryRead(() => + db + .select({ + id: externalServiceComponent.id, + externalServiceId: externalServiceComponent.externalServiceId, + upstreamComponentId: externalServiceComponent.upstreamComponentId, + slug: externalServiceComponent.slug, + aliases: externalServiceComponent.aliases, + name: externalServiceComponent.name, + description: externalServiceComponent.description, + groupName: externalServiceComponent.groupName, + position: externalServiceComponent.position, + indicator: externalServiceComponent.indicator, + status: externalServiceComponent.status, + firstSeenAt: externalServiceComponent.firstSeenAt, + }) + .from(externalServiceComponent) + .where( + and( + eq(externalServiceComponent.externalServiceId, service.id), + or( + eq(externalServiceComponent.slug, componentSlug), + sql`EXISTS (SELECT 1 FROM json_each(${externalServiceComponent.aliases}) WHERE value = ${componentSlug})`, + ), ), - ), - ) - .limit(1) - .all(); + ) + .limit(1) + .all(), + ); const row = rows[0]; if (!row) return { service, component: null }; diff --git a/packages/services/src/external-service-incident/list.ts b/packages/services/src/external-service-incident/list.ts index 1a1dcfdf..2d2c4baf 100644 --- a/packages/services/src/external-service-incident/list.ts +++ b/packages/services/src/external-service-incident/list.ts @@ -4,6 +4,7 @@ import { externalServiceIncident } from "@openstatus/db/src/schema"; import { getExternalServiceBySlug } from "../external-service"; import type { GlobalReadContext } from "../external-service/internal"; +import { retryRead } from "../retry"; export type ExternalIncidentListItem = { id: number; @@ -39,29 +40,31 @@ export async function listExternalIncidentsByServiceId(args: { const db = ctx?.db ?? defaultDb; const limit = args.limit ?? DEFAULT_LIMIT; - const rows = await db - .select({ - id: externalServiceIncident.id, - externalServiceId: externalServiceIncident.externalServiceId, - providerIncidentId: externalServiceIncident.providerIncidentId, - name: externalServiceIncident.name, - status: externalServiceIncident.status, - impact: externalServiceIncident.impact, - shortlink: externalServiceIncident.shortlink, - startedAt: externalServiceIncident.startedAt, - createdAt: externalServiceIncident.createdAt, - resolvedAt: externalServiceIncident.resolvedAt, - }) - .from(externalServiceIncident) - .where(eq(externalServiceIncident.externalServiceId, externalServiceId)) - .orderBy( - desc( - sql`COALESCE(${externalServiceIncident.startedAt}, ${externalServiceIncident.createdAt})`, - ), - desc(externalServiceIncident.createdAt), - ) - .limit(limit) - .all(); + const rows = await retryRead(() => + db + .select({ + id: externalServiceIncident.id, + externalServiceId: externalServiceIncident.externalServiceId, + providerIncidentId: externalServiceIncident.providerIncidentId, + name: externalServiceIncident.name, + status: externalServiceIncident.status, + impact: externalServiceIncident.impact, + shortlink: externalServiceIncident.shortlink, + startedAt: externalServiceIncident.startedAt, + createdAt: externalServiceIncident.createdAt, + resolvedAt: externalServiceIncident.resolvedAt, + }) + .from(externalServiceIncident) + .where(eq(externalServiceIncident.externalServiceId, externalServiceId)) + .orderBy( + desc( + sql`COALESCE(${externalServiceIncident.startedAt}, ${externalServiceIncident.createdAt})`, + ), + desc(externalServiceIncident.createdAt), + ) + .limit(limit) + .all(), + ); return rows; } @@ -78,34 +81,36 @@ export async function listExternalIncidentsByComponent(args: { const db = ctx?.db ?? defaultDb; const limit = args.limit ?? DEFAULT_LIMIT; - const rows = await db - .select({ - id: externalServiceIncident.id, - externalServiceId: externalServiceIncident.externalServiceId, - providerIncidentId: externalServiceIncident.providerIncidentId, - name: externalServiceIncident.name, - status: externalServiceIncident.status, - impact: externalServiceIncident.impact, - shortlink: externalServiceIncident.shortlink, - startedAt: externalServiceIncident.startedAt, - createdAt: externalServiceIncident.createdAt, - resolvedAt: externalServiceIncident.resolvedAt, - }) - .from(externalServiceIncident) - .where( - and( - eq(externalServiceIncident.externalServiceId, externalServiceId), - sql`EXISTS (SELECT 1 FROM json_each(${externalServiceIncident.affectedComponentIds}) WHERE value = ${upstreamComponentId})`, - ), - ) - .orderBy( - desc( - sql`COALESCE(${externalServiceIncident.startedAt}, ${externalServiceIncident.createdAt})`, - ), - desc(externalServiceIncident.createdAt), - ) - .limit(limit) - .all(); + const rows = await retryRead(() => + db + .select({ + id: externalServiceIncident.id, + externalServiceId: externalServiceIncident.externalServiceId, + providerIncidentId: externalServiceIncident.providerIncidentId, + name: externalServiceIncident.name, + status: externalServiceIncident.status, + impact: externalServiceIncident.impact, + shortlink: externalServiceIncident.shortlink, + startedAt: externalServiceIncident.startedAt, + createdAt: externalServiceIncident.createdAt, + resolvedAt: externalServiceIncident.resolvedAt, + }) + .from(externalServiceIncident) + .where( + and( + eq(externalServiceIncident.externalServiceId, externalServiceId), + sql`EXISTS (SELECT 1 FROM json_each(${externalServiceIncident.affectedComponentIds}) WHERE value = ${upstreamComponentId})`, + ), + ) + .orderBy( + desc( + sql`COALESCE(${externalServiceIncident.startedAt}, ${externalServiceIncident.createdAt})`, + ), + desc(externalServiceIncident.createdAt), + ) + .limit(limit) + .all(), + ); return rows; } diff --git a/packages/services/src/external-service-report/read.ts b/packages/services/src/external-service-report/read.ts index c1cfadf4..9a231a95 100644 --- a/packages/services/src/external-service-report/read.ts +++ b/packages/services/src/external-service-report/read.ts @@ -12,6 +12,7 @@ import { import { externalServiceReport } from "@openstatus/db/src/schema"; import type { GlobalReadContext } from "../external-service/internal"; +import { retryRead } from "../retry"; export type ServiceReportWindow = { externalServiceId: number; @@ -43,22 +44,24 @@ export async function getServiceReportWindows(args: { const { ctx, serviceIds, since } = args; if (serviceIds.length === 0) return []; const db = ctx?.db ?? defaultDb; - return db - .select({ - externalServiceId: externalServiceReport.externalServiceId, - reporters: distinctReporters, - total: totalReports, - countries: distinctCountries, - }) - .from(externalServiceReport) - .where( - and( - inArray(externalServiceReport.externalServiceId, serviceIds), - gte(externalServiceReport.createdAt, since), - ), - ) - .groupBy(externalServiceReport.externalServiceId) - .all(); + return retryRead(() => + db + .select({ + externalServiceId: externalServiceReport.externalServiceId, + reporters: distinctReporters, + total: totalReports, + countries: distinctCountries, + }) + .from(externalServiceReport) + .where( + and( + inArray(externalServiceReport.externalServiceId, serviceIds), + gte(externalServiceReport.createdAt, since), + ), + ) + .groupBy(externalServiceReport.externalServiceId) + .all(), + ); } export async function getComponentReportWindows(args: { @@ -68,22 +71,24 @@ export async function getComponentReportWindows(args: { }): Promise { const { ctx, serviceId, since } = args; const db = ctx?.db ?? defaultDb; - const rows = await db - .select({ - componentId: externalServiceReport.externalServiceComponentId, - reporters: distinctReporters, - total: totalReports, - }) - .from(externalServiceReport) - .where( - and( - eq(externalServiceReport.externalServiceId, serviceId), - isNotNull(externalServiceReport.externalServiceComponentId), - gte(externalServiceReport.createdAt, since), - ), - ) - .groupBy(externalServiceReport.externalServiceComponentId) - .all(); + const rows = await retryRead(() => + db + .select({ + componentId: externalServiceReport.externalServiceComponentId, + reporters: distinctReporters, + total: totalReports, + }) + .from(externalServiceReport) + .where( + and( + eq(externalServiceReport.externalServiceId, serviceId), + isNotNull(externalServiceReport.externalServiceComponentId), + gte(externalServiceReport.createdAt, since), + ), + ) + .groupBy(externalServiceReport.externalServiceComponentId) + .all(), + ); return rows.flatMap((r) => r.componentId == null ? [] @@ -104,22 +109,24 @@ export async function getServiceReportDaily(args: { }): Promise { const { ctx, serviceId, since } = args; const db = ctx?.db ?? defaultDb; - return db - .select({ - day: dayBucket, - reporters: distinctReporters, - total: totalReports, - }) - .from(externalServiceReport) - .where( - and( - eq(externalServiceReport.externalServiceId, serviceId), - gte(externalServiceReport.createdAt, since), - ), - ) - .groupBy(dayBucket) - .orderBy(asc(dayBucket)) - .all(); + return retryRead(() => + db + .select({ + day: dayBucket, + reporters: distinctReporters, + total: totalReports, + }) + .from(externalServiceReport) + .where( + and( + eq(externalServiceReport.externalServiceId, serviceId), + gte(externalServiceReport.createdAt, since), + ), + ) + .groupBy(dayBucket) + .orderBy(asc(dayBucket)) + .all(), + ); } export async function getServiceReportCountries(args: { @@ -131,18 +138,20 @@ export async function getServiceReportCountries(args: { const { ctx, serviceId, since } = args; const limit = Math.max(0, Math.trunc(args.limit)); const db = ctx?.db ?? defaultDb; - return db - .select({ country: externalServiceReport.country, total: totalReports }) - .from(externalServiceReport) - .where( - and( - eq(externalServiceReport.externalServiceId, serviceId), - gte(externalServiceReport.createdAt, since), - sql`${externalServiceReport.country} != ''`, - ), - ) - .groupBy(externalServiceReport.country) - .orderBy(desc(totalReports), asc(externalServiceReport.country)) - .limit(limit) - .all(); + return retryRead(() => + db + .select({ country: externalServiceReport.country, total: totalReports }) + .from(externalServiceReport) + .where( + and( + eq(externalServiceReport.externalServiceId, serviceId), + gte(externalServiceReport.createdAt, since), + sql`${externalServiceReport.country} != ''`, + ), + ) + .groupBy(externalServiceReport.country) + .orderBy(desc(totalReports), asc(externalServiceReport.country)) + .limit(limit) + .all(), + ); } diff --git a/packages/services/src/external-service/list-slugs.ts b/packages/services/src/external-service/list-slugs.ts index 8d098fce..2f886af3 100644 --- a/packages/services/src/external-service/list-slugs.ts +++ b/packages/services/src/external-service/list-slugs.ts @@ -1,6 +1,7 @@ import { db as defaultDb } from "@openstatus/db"; import { externalService } from "@openstatus/db/src/schema"; +import { retryRead } from "../retry"; import { type GlobalReadContext, liveOnlyClause } from "./internal"; export type SlugMap = { @@ -22,9 +23,9 @@ export async function listExternalServiceSlugs(args?: { }) .from(externalService); - const rows = includeDeleted - ? await query.all() - : await query.where(liveOnlyClause()).all(); + const rows = await retryRead(() => + includeDeleted ? query.all() : query.where(liveOnlyClause()).all(), + ); const canonical: string[] = []; const aliases: Array<{ from: string; to: string }> = []; diff --git a/packages/services/src/external-service/list.ts b/packages/services/src/external-service/list.ts index 0707e5a0..c539d9e8 100644 --- a/packages/services/src/external-service/list.ts +++ b/packages/services/src/external-service/list.ts @@ -6,6 +6,7 @@ import { selectExternalServiceSchema, } from "@openstatus/db/src/schema"; +import { retryRead } from "../retry"; import { type GlobalReadContext, liveOnlyClause } from "./internal"; export type ListExternalServicesInput = { @@ -57,9 +58,9 @@ export async function listExternalServices(args: { .from(externalService) .orderBy(asc(sql`lower(${externalService.name})`)); - const rows = input?.includeDeleted - ? await query.all() - : await query.where(liveOnlyClause()).all(); + const rows = await retryRead(() => + input?.includeDeleted ? query.all() : query.where(liveOnlyClause()).all(), + ); return validateRows(rows); } @@ -71,16 +72,18 @@ export async function getExternalServiceBySlug(args: { const { ctx, slug } = args; const db = ctx?.db ?? defaultDb; - const rows = await db - .select() - .from(externalService) - .where( - or( - eq(externalService.slug, slug), - sql`EXISTS (SELECT 1 FROM json_each(${externalService.aliases}) WHERE value = ${slug})`, - ), - ) - .all(); + const rows = await retryRead(() => + db + .select() + .from(externalService) + .where( + or( + eq(externalService.slug, slug), + sql`EXISTS (SELECT 1 FROM json_each(${externalService.aliases}) WHERE value = ${slug})`, + ), + ) + .all(), + ); const validated = validateRows(rows); return validated[0] ?? null; diff --git a/packages/services/src/index.ts b/packages/services/src/index.ts index f68adba0..7b10554f 100644 --- a/packages/services/src/index.ts +++ b/packages/services/src/index.ts @@ -14,6 +14,7 @@ export { export { isRetryableDbError, isTransientServerError, + retryRead, withBusyRetry, } from "./retry"; diff --git a/packages/services/src/retry.ts b/packages/services/src/retry.ts index 98bd7dee..1e649f78 100644 --- a/packages/services/src/retry.ts +++ b/packages/services/src/retry.ts @@ -1,9 +1,10 @@ +import { Cause, Effect, Exit, Schedule } from "effect"; + const RETRYABLE_CODES = new Set(["SQLITE_BUSY", "SQLITE_LOCKED"]); const RETRYABLE_MESSAGE = /database is (locked|busy)/i; const TRANSIENT_SERVER_MESSAGE = /Server returned HTTP status 5\d\d/i; const MAX_ATTEMPTS = 5; const BASE_DELAY_MS = 25; -const MAX_DELAY_MS = 400; const MAX_CAUSE_DEPTH = 10; function hasCode(value: object): value is { code: unknown } { @@ -61,23 +62,28 @@ export function isTransientServerError(err: unknown): boolean { return false; } -function backoffDelayMs(attempt: number): number { - const cap = Math.min(MAX_DELAY_MS, BASE_DELAY_MS * 2 ** attempt); - return Math.random() * cap; -} - -const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); +const retrySchedule = Schedule.exponential(`${BASE_DELAY_MS} millis`).pipe( + Schedule.jittered, +); export async function withBusyRetry( fn: () => Promise, isRetryable: (err: unknown) => boolean = isRetryableDbError, ): Promise { - for (let attempt = 0; ; attempt++) { - try { - return await fn(); - } catch (err) { - if (attempt >= MAX_ATTEMPTS - 1 || !isRetryable(err)) throw err; - await sleep(backoffDelayMs(attempt)); - } - } + const exit = await Effect.runPromiseExit( + Effect.tryPromise({ try: () => fn(), catch: (err) => err }).pipe( + Effect.retry({ + schedule: retrySchedule, + times: MAX_ATTEMPTS - 1, + while: isRetryable, + }), + ), + ); + // squash so callers reject with the original error, not Effect's FiberFailure. + if (Exit.isFailure(exit)) throw Cause.squash(exit.cause); + return exit.value; } + +// Idempotent reads only: also retries transient Turso 5xx (see isTransientServerError). +export const retryRead = (fn: () => Promise): Promise => + withBusyRetry(fn, (e) => isRetryableDbError(e) || isTransientServerError(e)); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index c0a3f9b6..a7555312 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -2551,6 +2551,9 @@ importers: '@openstatus/utils': specifier: workspace:* version: link:../utils + effect: + specifier: 'catalog:' + version: 3.21.2 zod: specifier: 'catalog:' version: 4.1.13 -- 2.51.2