From 963902514b737f9163dd6ddecb203e30e277f5d9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ullrich=20Sch=C3=A4fer?= Date: Tue, 21 Apr 2026 07:36:10 +0200 Subject: [PATCH] Failover across multiple Overpass upstreams MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Public Overpass instances go bad without warning — we just lived through a day where overpass.private.coffee (our single upstream) returned 504s after 60s while lz4.overpass-api.de served the same query in 0.6s. Survive that by trying a list of upstreams in order until one responds usefully. Server changes: - OVERPASS_URLS env — comma-separated list, tried in order. Falls back to OVERPASS_URL for backward compat, then to a built-in default pair (lz4.overpass-api.de, overpass-api.de). - Per-upstream timeout of 10s via AbortSignal.timeout, so a saturated instance fails over in seconds, not minutes. - Client abort propagation via AbortSignal.any — if the browser cancels (e.g. user panned the map), the in-flight upstream fetch is aborted too. Stops wasting upstream capacity on requests nobody wants. - Rate-limited-body detection (Overpass returns 200 with a `rate_limited` marker when throttling) triggers failover. - Per-upstream labels on overpass_upstream_requests_total and overpass_upstream_duration_seconds so the dashboard can break out health + latency by instance. Histogram buckets extended to 30s. Compose: OVERPASS_URLS passthrough with the default pair hard-wired, overridable via SOPS. Tests: 8 cases covering the new fetchWithFailover helper — happy path, 5xx failover, rate_limited failover, network error failover, timeout tagging, all-upstreams-fail, client abort stops the loop, per-attempt latency observation. Follow-ups left out of scope: - Client-side debounce (UI concern). - Moving `.observe()` after body read for accuracy in measurement. Co-Authored-By: Claude Opus 4.7 (1M context) --- apps/planner/app/lib/metrics.server.ts | 11 +- apps/planner/app/routes/api.overpass.test.ts | 209 +++++++++++++++++++ apps/planner/app/routes/api.overpass.ts | 134 ++++++++++-- apps/planner/vitest.config.ts | 17 +- infrastructure/.env.example | 5 + infrastructure/docker-compose.yml | 7 +- 6 files changed, 358 insertions(+), 25 deletions(-) create mode 100644 apps/planner/app/routes/api.overpass.test.ts diff --git a/apps/planner/app/lib/metrics.server.ts b/apps/planner/app/lib/metrics.server.ts index b89f8c2..2aa8d9d 100644 --- a/apps/planner/app/lib/metrics.server.ts +++ b/apps/planner/app/lib/metrics.server.ts @@ -60,15 +60,20 @@ export const overpassUpstreamDuration = getOrCreate("overpass_upstream_duration_ new client.Histogram({ name: "overpass_upstream_duration_seconds", help: "Duration of upstream Overpass API requests in seconds", - buckets: [0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10], + labelNames: ["upstream"] as const, + // Buckets go to 30s because our per-upstream timeout is 10s; a bit + // of headroom catches slow-but-not-timed-out tails without overly + // coarse resolution at the fast end where lz4 typically lands + // (~100–500ms). + buckets: [0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 15, 30], }), ); export const overpassUpstreamRequests = getOrCreate("overpass_upstream_requests_total", () => new client.Counter({ name: "overpass_upstream_requests_total", - help: "Upstream Overpass API requests by status", - labelNames: ["status"] as const, + help: "Upstream Overpass API requests by upstream host and status", + labelNames: ["upstream", "status"] as const, }), ); diff --git a/apps/planner/app/routes/api.overpass.test.ts b/apps/planner/app/routes/api.overpass.test.ts new file mode 100644 index 0000000..2d4a6e0 --- /dev/null +++ b/apps/planner/app/routes/api.overpass.test.ts @@ -0,0 +1,209 @@ +import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; + +// Metrics module pulls in prom-client and registers processes-wide +// metrics on import; mock it out so tests aren't entangled with +// the registry and can assert on label values instead. +vi.mock("~/lib/metrics.server", () => ({ + overpassCacheEvents: { inc: vi.fn() }, + overpassCacheSize: { set: vi.fn() }, + overpassUpstreamDuration: { observe: vi.fn() }, + overpassUpstreamRequests: { inc: vi.fn() }, +})); + +import { fetchWithFailover } from "./api.overpass"; +import { + overpassUpstreamRequests, + overpassUpstreamDuration, +} from "~/lib/metrics.server"; + +const URL_A = "https://a.example/api/interpreter"; +const URL_B = "https://b.example/api/interpreter"; +const URL_C = "https://c.example/api/interpreter"; + +function makeResponse( + body: string, + init: { status?: number; headers?: Record } = {}, +): Response { + return new Response(body, { + status: init.status ?? 200, + headers: init.headers ?? { "content-type": "application/json" }, + }); +} + +beforeEach(() => { + // `clearAllMocks` zeroes call-history on module-level mocks like + // `overpassUpstreamRequests.inc`, which would otherwise accumulate + // across test cases and break assertions that inspect call counts. + vi.clearAllMocks(); +}); + +afterEach(() => { + vi.unstubAllGlobals(); +}); + +describe("fetchWithFailover", () => { + it("returns the first upstream's response when it succeeds", async () => { + const fetchSpy = vi.fn().mockResolvedValueOnce(makeResponse('{"ok":1}')); + vi.stubGlobal("fetch", fetchSpy); + + const result = await fetchWithFailover( + "data=...", + new AbortController().signal, + [URL_A, URL_B], + ); + + expect(result?.status).toBe(200); + expect(result?.body).toBe('{"ok":1}'); + expect(fetchSpy).toHaveBeenCalledTimes(1); + expect(fetchSpy).toHaveBeenCalledWith(URL_A, expect.any(Object)); + expect(overpassUpstreamRequests.inc).toHaveBeenCalledWith({ + upstream: "a.example", + status: "200", + }); + }); + + it("falls over to the next upstream on non-2xx", async () => { + const fetchSpy = vi + .fn() + .mockResolvedValueOnce(makeResponse("gateway timeout", { status: 504 })) + .mockResolvedValueOnce(makeResponse('{"elements":[]}')); + vi.stubGlobal("fetch", fetchSpy); + + const result = await fetchWithFailover( + "data=...", + new AbortController().signal, + [URL_A, URL_B], + ); + + expect(result?.body).toBe('{"elements":[]}'); + expect(fetchSpy).toHaveBeenCalledTimes(2); + expect(overpassUpstreamRequests.inc).toHaveBeenCalledWith({ + upstream: "a.example", + status: "504", + }); + expect(overpassUpstreamRequests.inc).toHaveBeenCalledWith({ + upstream: "b.example", + status: "200", + }); + }); + + it("falls over when the body carries Overpass's rate_limited runtime error", async () => { + // Overpass returns HTTP 200 with a `rate_limited` marker in the + // body when it's throttling. Treat that as an upstream failure. + const fetchSpy = vi + .fn() + .mockResolvedValueOnce(makeResponse('{"error":"rate_limited"}')) + .mockResolvedValueOnce(makeResponse('{"elements":[]}')); + vi.stubGlobal("fetch", fetchSpy); + + const result = await fetchWithFailover( + "data=...", + new AbortController().signal, + [URL_A, URL_B], + ); + + expect(result?.body).toBe('{"elements":[]}'); + expect(fetchSpy).toHaveBeenCalledTimes(2); + }); + + it("falls over on network / timeout errors", async () => { + const fetchSpy = vi + .fn() + .mockRejectedValueOnce(new Error("connect ECONNREFUSED")) + .mockResolvedValueOnce(makeResponse('{"elements":[]}')); + vi.stubGlobal("fetch", fetchSpy); + + const result = await fetchWithFailover( + "data=...", + new AbortController().signal, + [URL_A, URL_B], + ); + + expect(result?.body).toBe('{"elements":[]}'); + expect(fetchSpy).toHaveBeenCalledTimes(2); + expect(overpassUpstreamRequests.inc).toHaveBeenCalledWith({ + upstream: "a.example", + status: "error", + }); + }); + + it("tags TimeoutError distinctly from other errors", async () => { + const timeoutErr = new Error("aborted"); + timeoutErr.name = "TimeoutError"; + const fetchSpy = vi + .fn() + .mockRejectedValueOnce(timeoutErr) + .mockResolvedValueOnce(makeResponse('{"elements":[]}')); + vi.stubGlobal("fetch", fetchSpy); + + await fetchWithFailover( + "data=...", + new AbortController().signal, + [URL_A, URL_B], + ); + + expect(overpassUpstreamRequests.inc).toHaveBeenCalledWith({ + upstream: "a.example", + status: "timeout", + }); + }); + + it("returns null when every upstream fails", async () => { + const fetchSpy = vi + .fn() + .mockResolvedValueOnce(makeResponse("bad", { status: 500 })) + .mockResolvedValueOnce(makeResponse("bad", { status: 502 })) + .mockResolvedValueOnce(makeResponse("bad", { status: 503 })); + vi.stubGlobal("fetch", fetchSpy); + + const result = await fetchWithFailover( + "data=...", + new AbortController().signal, + [URL_A, URL_B, URL_C], + ); + + expect(result).toBeNull(); + expect(fetchSpy).toHaveBeenCalledTimes(3); + }); + + it("stops iterating when the client aborts, rethrows, and doesn't try later upstreams", async () => { + const controller = new AbortController(); + const abortErr = new DOMException("aborted", "AbortError"); + const fetchSpy = vi.fn().mockImplementationOnce(async () => { + controller.abort(); + throw abortErr; + }); + vi.stubGlobal("fetch", fetchSpy); + + await expect( + fetchWithFailover("data=...", controller.signal, [URL_A, URL_B]), + ).rejects.toBe(abortErr); + + // Only A was attempted; B never got a chance. + expect(fetchSpy).toHaveBeenCalledTimes(1); + expect(overpassUpstreamRequests.inc).toHaveBeenCalledWith({ + upstream: "a.example", + status: "client-abort", + }); + }); + + it("records latency with the upstream label on every attempt (success and failure)", async () => { + const fetchSpy = vi + .fn() + .mockResolvedValueOnce(makeResponse("bad", { status: 503 })) + .mockResolvedValueOnce(makeResponse('{"elements":[]}')); + vi.stubGlobal("fetch", fetchSpy); + + await fetchWithFailover( + "data=...", + new AbortController().signal, + [URL_A, URL_B], + ); + + const calls = vi.mocked(overpassUpstreamDuration.observe).mock.calls; + const upstreams = calls.map( + (c) => (c[0] as unknown as { upstream: string }).upstream, + ); + expect(upstreams).toEqual(["a.example", "b.example"]); + }); +}); diff --git a/apps/planner/app/routes/api.overpass.ts b/apps/planner/app/routes/api.overpass.ts index 571e97a..b1cb400 100644 --- a/apps/planner/app/routes/api.overpass.ts +++ b/apps/planner/app/routes/api.overpass.ts @@ -7,9 +7,33 @@ import { overpassUpstreamRequests, } from "~/lib/metrics.server"; -const UPSTREAM_URL = process.env.OVERPASS_URL ?? "https://overpass.private.coffee/api/interpreter"; +/** + * Ordered list of upstream Overpass endpoints. We try each in turn + * until one returns a usable response; on timeout, network error, + * non-2xx status, or a body carrying Overpass's `rate_limited` runtime + * error, we fall over to the next. The default list keeps us on + * healthy community instances; `OVERPASS_URLS` overrides it and + * `OVERPASS_URL` is kept as a single-entry backward-compat alias. + */ +const DEFAULT_UPSTREAMS = [ + "https://lz4.overpass-api.de/api/interpreter", + "https://overpass-api.de/api/interpreter", +]; + +function loadUpstreams(): string[] { + const list = process.env.OVERPASS_URLS ?? process.env.OVERPASS_URL; + if (!list) return DEFAULT_UPSTREAMS; + return list.split(",").map((s) => s.trim()).filter(Boolean); +} + +const UPSTREAMS = loadUpstreams(); const USER_AGENT = "trails.cool Planner (https://trails.cool; legal@trails.cool)"; +// Per-attempt wall-clock budget. Keep aggressive so failover kicks in +// quickly when an upstream is saturated — a single slow instance +// shouldn't cost the user minutes. +const PER_UPSTREAM_TIMEOUT_MS = 10_000; + const CACHE_TTL_MS = 10 * 60 * 1000; const CACHE_MAX_ENTRIES = 200; @@ -19,8 +43,15 @@ interface CacheEntry { expiresAt: number; } +interface UpstreamResult { + body: string; + contentType: string; + status: number; + cacheable: boolean; +} + const cache = new Map(); -const inFlight = new Map>(); +const inFlight = new Map>(); function getFromCache(key: string): CacheEntry | null { const entry = cache.get(key); @@ -45,6 +76,80 @@ function putInCache(key: string, body: string, contentType: string) { overpassCacheSize.set(cache.size); } +function hostOf(url: string): string { + try { + return new URL(url).hostname; + } catch { + return "unknown"; + } +} + +/** + * Try each upstream in order until one returns a usable response. + * Returns the first successful result, or `null` if every upstream + * failed. `clientSignal` is the browser's abort signal — if the user + * gives up before we succeed, we stop trying and propagate the abort + * to the in-flight upstream fetch so we don't waste its capacity on a + * request nobody wants anymore. Exported for tests; the action + * handler calls it with the module-level `UPSTREAMS` list. + */ +export async function fetchWithFailover( + body: string, + clientSignal: AbortSignal, + urls: readonly string[] = UPSTREAMS, +): Promise { + for (const url of urls) { + if (clientSignal.aborted) break; + const upstream = hostOf(url); + const start = Date.now(); + try { + const signal = AbortSignal.any([ + clientSignal, + AbortSignal.timeout(PER_UPSTREAM_TIMEOUT_MS), + ]); + const response = await fetch(url, { + method: "POST", + headers: { + "Content-Type": "application/x-www-form-urlencoded", + "User-Agent": USER_AGENT, + }, + body, + signal, + }); + const responseBody = await response.text(); + overpassUpstreamDuration.observe({ upstream }, (Date.now() - start) / 1000); + overpassUpstreamRequests.inc({ upstream, status: String(response.status) }); + + // Overpass returns 200 with a `rate_limited` runtime error when + // the instance is throttling us. Treat that like a 5xx: try the + // next upstream. + if (!response.ok || responseBody.includes("rate_limited")) { + continue; + } + + return { + body: responseBody, + contentType: response.headers.get("content-type") ?? "application/json", + status: 200, + cacheable: true, + }; + } catch (err) { + overpassUpstreamDuration.observe({ upstream }, (Date.now() - start) / 1000); + // If the CLIENT aborted, stop the whole loop — the user doesn't + // want an answer anymore. Any other error (timeout, network, + // DNS) counts as "this upstream failed, try the next." + if (clientSignal.aborted) { + overpassUpstreamRequests.inc({ upstream, status: "client-abort" }); + throw err; + } + const label = (err as Error).name === "TimeoutError" ? "timeout" : "error"; + overpassUpstreamRequests.inc({ upstream, status: label }); + continue; + } + } + return null; +} + export async function action({ request }: Route.ActionArgs) { if (request.method !== "POST") { return new Response("Method not allowed", { status: 405 }); @@ -91,23 +196,13 @@ export async function action({ request }: Route.ActionArgs) { let pending = inFlight.get(cacheKey); if (!pending) { overpassCacheEvents.inc({ result: "miss" }); - pending = (async () => { - const start = Date.now(); + pending = (async (): Promise => { try { - const upstream = await fetch(UPSTREAM_URL, { - method: "POST", - headers: { - "Content-Type": "application/x-www-form-urlencoded", - "User-Agent": USER_AGENT, - }, - body, - }); - overpassUpstreamDuration.observe((Date.now() - start) / 1000); - overpassUpstreamRequests.inc({ status: String(upstream.status) }); - const responseBody = await upstream.text(); - const contentType = upstream.headers.get("content-type") ?? "application/json"; - const cacheable = upstream.status === 200 && !responseBody.includes("rate_limited"); - return { body: responseBody, contentType, status: upstream.status, cacheable }; + const result = await fetchWithFailover(body, request.signal); + if (!result) { + throw new Error("all upstreams failed"); + } + return result; } finally { inFlight.delete(cacheKey); } @@ -117,11 +212,10 @@ export async function action({ request }: Route.ActionArgs) { overpassCacheEvents.inc({ result: "coalesced" }); } - let result; + let result: UpstreamResult; try { result = await pending; } catch { - overpassUpstreamRequests.inc({ status: "error" }); return new Response("Upstream Overpass unavailable", { status: 502 }); } diff --git a/apps/planner/vitest.config.ts b/apps/planner/vitest.config.ts index 8a07b14..cf58c77 100644 --- a/apps/planner/vitest.config.ts +++ b/apps/planner/vitest.config.ts @@ -1 +1,16 @@ -export { default } from "../../vitest.shared.ts"; +import { defineConfig, mergeConfig } from "vitest/config"; +import shared from "../../vitest.shared.ts"; +import { resolve } from "node:path"; + +// Mirror the `~` alias that React Router's runtime provides so test +// files can resolve `~/lib/...` the same way the route modules do. +export default mergeConfig( + shared, + defineConfig({ + resolve: { + alias: { + "~": resolve(import.meta.dirname, "app"), + }, + }, + }), +); diff --git a/infrastructure/.env.example b/infrastructure/.env.example index ef8cd88..85d59dc 100644 --- a/infrastructure/.env.example +++ b/infrastructure/.env.example @@ -7,6 +7,11 @@ JWT_SECRET=change-me S3_ACCESS_KEY= S3_SECRET_KEY= +# Overpass upstream failover. Comma-separated list tried in order; +# first 2xx response wins. Falls back to the single-URL OVERPASS_URL +# if this is unset, and to a built-in default list if neither is set. +# OVERPASS_URLS=https://lz4.overpass-api.de/api/interpreter,https://overpass-api.de/api/interpreter + # Demo-activity-bot. Only enable in prod. When true, the journal worker # bootstraps a demo user and generates public demo routes + activities # every 90 min. Retention is in days (default 14). diff --git a/infrastructure/docker-compose.yml b/infrastructure/docker-compose.yml index 7010104..e4d80d5 100644 --- a/infrastructure/docker-compose.yml +++ b/infrastructure/docker-compose.yml @@ -55,7 +55,12 @@ services: restart: unless-stopped environment: BROUTER_URL: http://brouter:17777 - OVERPASS_URL: https://overpass.private.coffee/api/interpreter + # Ordered failover list for the Overpass proxy. The code defaults + # to the same pair if unset; declaring it here makes the prod + # upstream explicit and easy to reshuffle via env when a + # community instance misbehaves. Override per-env with + # OVERPASS_URLS in the SOPS file. + OVERPASS_URLS: ${OVERPASS_URLS:-https://lz4.overpass-api.de/api/interpreter,https://overpass-api.de/api/interpreter} DATABASE_URL: postgres://trails:${POSTGRES_PASSWORD:-trails}@postgres:5432/trails NODE_ENV: production PORT: 3001