diff --git a/docs/api-logs.md b/docs/api-logs.md new file mode 100644 index 00000000..4b2edb7a --- /dev/null +++ b/docs/api-logs.md @@ -0,0 +1,52 @@ +# API Logs Proxy Endpoint + +The `/api/logs` endpoint acts as a reverse proxy for downstream logging services, forwarding requests to the configured `UPSTREAM_URL/logs` and applying per-endpoint circuit breaker protection to prevent cascading failures. + +## Features + +- **Per-Endpoint Circuit Breaking**: Each distinct downstream endpoint (e.g. `/api/logs/system` vs `/api/logs/audit`) is tracked independently by the `BreakerRegistry`. +- **Fast-Fail on Open**: If the downstream service begins failing and the circuit breaker trips to `OPEN`, the gateway fast-fails immediately with HTTP `503 Service Unavailable`, protecting both the gateway resources and the downstream logging service during an outage. +- **Support for All Methods**: The endpoint supports arbitrary HTTP verbs (`GET`, `POST`, `PUT`, `DELETE`, etc.) and seamlessly proxies the request body. + +## API Reference + +### `ALL /api/logs/:endpoint(*)` + +Proxies the request to the upstream logging server. + +**Path Parameters**: +- `endpoint` (optional): The specific downstream log path. For example, `GET /api/logs/system/metrics` will proxy to `UPSTREAM_URL/logs/system/metrics`. + +**Headers**: +- Passes through `Authorization` and `Content-Type`. + +## Circuit Breaker Behavior + +The underlying circuit breaker is configured with the following defaults: +- **Failure Threshold**: 5 consecutive failures before opening. +- **Cooldown**: 30 seconds before attempting a half-open probe. +- **Success Threshold**: 1 successful probe to close the breaker and resume normal traffic. + +When the breaker is `OPEN`, the gateway returns: + +```json +{ + "success": false, + "error": { + "code": "SERVICE_UNAVAILABLE", + "message": "Downstream logs endpoint is currently unavailable: Circuit breaker is open. Cooldown remaining: 29999ms" + } +} +``` + +Unexpected errors (e.g. invalid response format, connection timeout) are returned as: + +```json +{ + "success": false, + "error": { + "code": "BAD_GATEWAY", + "message": "Upstream error message" + } +} +``` diff --git a/src/index.ts b/src/index.ts index e9aaad96..71dd7779 100644 --- a/src/index.ts +++ b/src/index.ts @@ -24,6 +24,38 @@ import { } from "./lifecycle/shutdown.js"; import type { Socket } from "net"; +import { createDeveloperRouter } from "./routes/developerRoutes.js"; +import { createGatewayRouter } from "./routes/gatewayRoutes.js"; +import { createProxyRouter } from "./routes/proxyRoutes.js"; +import adminRouter from "./routes/admin.js"; +import logsRouter from "./routes/logs.js"; +import { createUsageAnomaliesRouter } from "./routes/admin/usage/anomalies.js"; +import refundsRouter from "./routes/refunds.js"; +import { defaultDeveloperRepository } from "./repositories/developerRepository.js"; +import { createBillingService } from "./services/billingService.js"; +import { + createConfiguredRateLimiter, + resolveRateLimiterConfig, +} from "./services/rateLimiter.js"; +import { PgUsageEventsRepository } from "./repositories/usageEventsRepository.pg.js"; +import { createRevenueLedgerIndexerJob } from "./services/revenueLedgerIndexer.js"; +import { RevenueSettlementService } from "./services/revenueSettlementService.js"; +import { createSettlementStatusSyncJob } from "./services/settlementStatusSyncJob.js"; +import { createIdempotencySweeperJob } from "./services/idempotencySweeper.js"; +import { createPostgresUsageStore } from "./services/usageStore.js"; +import { createPostgresSettlementStore } from "./services/settlementStore.js"; +import { createApiRegistry } from "./data/apiRegistry.js"; +import { ApiKey } from "./types/gateway.js"; +import { listingsCache } from "./lib/listingsCache.js"; +import { createSlowQueryAlerterJob } from "./workers/slowQueryAlerter.js"; +import { createAnomalyDetectorJob } from "./workers/anomalyDetector.js"; +import { + initSloRecorder, + sloRecorderMiddleware, +} from "./workers/sloAlertRecorder.js"; +import { createSloAlertJob } from "./workers/sloAlertJob.js"; +import { createMonthlyInvoiceJob } from "./workers/monthlyInvoiceJob.js"; +import { createSettlementReconWorker } from "./workers/settlementRecon.js"; import { createDeveloperRouter } from './routes/developerRoutes.js'; import { createGatewayRouter } from './routes/gatewayRoutes.js'; import { createProxyRouter } from './routes/proxyRoutes.js'; @@ -231,6 +263,10 @@ if (isDirectExecution) { app.use("/api/developers", developerRouter); // Mounted before the generic admin router so it is not shadowed by // adminRouter's `/usage/:developerId` route. + app.use("/api/admin/usage/anomalies", createUsageAnomaliesRouter({ pool })); + app.use("/api/admin", adminRouter); + app.use("/api/refunds", refundsRouter); + app.use("/api/logs", logsRouter); app.use('/api/admin/usage/anomalies', createUsageAnomaliesRouter({ pool })); // Webhook management routes diff --git a/src/lib/circuitBreaker.ts b/src/lib/circuitBreaker.ts index 96ba02e3..a56d8ce3 100644 --- a/src/lib/circuitBreaker.ts +++ b/src/lib/circuitBreaker.ts @@ -28,6 +28,7 @@ export interface CircuitBreakerConfig { failureThreshold?: number; cooldownMs?: number; successThreshold?: number; + onOpenError?: (message: string) => Error; } export interface CircuitBreakerMetrics { @@ -299,20 +300,16 @@ export class CircuitBreaker { metrics = this.transitionTo(metrics, CircuitBreakerState.HALF_OPEN, now, breakerKey); await this.store.set(breakerKey, metrics); } else { - throw new CircuitBreakerOpenError( - `Circuit breaker is open. Cooldown remaining: ${ - this.config.cooldownMs - timeSinceFailure - }ms` - ); + const msg = `Circuit breaker is open. Cooldown remaining: ${this.config.cooldownMs - timeSinceFailure}ms`; + throw this.config.onOpenError ? this.config.onOpenError(msg) : new CircuitBreakerOpenError(msg); } } // Enforce exactly one trial call in HALF_OPEN state if (metrics.state === CircuitBreakerState.HALF_OPEN) { if (this.activeTrials.has(breakerKey)) { - throw new CircuitBreakerOpenError( - 'Circuit breaker is in half-open state. Only one trial call allowed at a time.' - ); + const msg = 'Circuit breaker is in half-open state. Only one trial call allowed at a time.'; + throw this.config.onOpenError ? this.config.onOpenError(msg) : new CircuitBreakerOpenError(msg); } this.activeTrials.add(breakerKey); } diff --git a/src/routes/logs.test.ts b/src/routes/logs.test.ts new file mode 100644 index 00000000..24dc47f4 --- /dev/null +++ b/src/routes/logs.test.ts @@ -0,0 +1,58 @@ +import express from "express"; +import request from "supertest"; +import logsRouter from "./logs.js"; +import { getDefaultBreakerRegistry, CircuitBreakerState } from "../lib/circuitBreaker.js"; +import { config } from "../config/index.js"; +import { errorHandler } from "../middleware/errorHandler.js"; + +const app = express(); +app.use(express.json()); +app.use("/api/logs", logsRouter); +app.use(errorHandler); + +const originalUpstreamUrl = config.proxy.upstreamUrl; + +describe("GET /api/logs", () => { + const breakerRegistry = getDefaultBreakerRegistry(); + + beforeEach(async () => { + // Reset all breakers before each test + const breakers = await breakerRegistry.list(); + for (const b of breakers) { + const breaker = breakerRegistry.get(b.slug); + if (breaker) { + await breaker.reset(b.slug); + } + } + }); + + afterAll(() => { + // Restore config + Object.assign(config.proxy, { upstreamUrl: originalUpstreamUrl }); + }); + + it("returns 503 fast on open circuit breaker", async () => { + // Point upstream to an invalid/failing URL to trigger circuit breaker failures + Object.assign(config.proxy, { upstreamUrl: "http://localhost:1" }); + + // The threshold in the logs route is 5. Let's make 5 failing requests + for (let i = 0; i < 5; i++) { + await request(app).get("/api/logs/test-endpoint").expect(502); + } + + // Circuit should now be OPEN, making it fast-fail with 503 + const response = await request(app).get("/api/logs/test-endpoint").expect(503); + + expect(response.body).toMatchObject({ + success: false, + error: { + code: "SERVICE_UNAVAILABLE", + }, + }); + expect(response.body.error.message).toMatch(/unavailable/i); + + // Verify breaker state is OPEN + const state = await breakerRegistry.getState('logs-get-test-endpoint'); + expect(state).toBe(CircuitBreakerState.OPEN); + }); +}); diff --git a/src/routes/logs.ts b/src/routes/logs.ts new file mode 100644 index 00000000..a3947fab --- /dev/null +++ b/src/routes/logs.ts @@ -0,0 +1,68 @@ +import { Router, Request, Response, NextFunction } from 'express'; +import { getDefaultBreakerRegistry } from '../lib/circuitBreaker.js'; +import { ServiceUnavailableError, BadGatewayError } from '../errors/index.js'; +import { config } from '../config/index.js'; + +const logsRouter = Router(); + +logsRouter.all('/:endpoint(*)?', async (req: Request, res: Response, next: NextFunction) => { + const endpoint = req.params.endpoint || 'default'; + const breakerKey = `logs-${req.method.toLowerCase()}-${endpoint.replace(/\//g, '-')}`; + + const breaker = getDefaultBreakerRegistry().getOrCreate(breakerKey, { + failureThreshold: 5, + cooldownMs: 30000, + successThreshold: 1, + onOpenError: (msg) => new ServiceUnavailableError(`Downstream logs endpoint is currently unavailable: ${msg}`), + }); + + try { + const result = await breaker.execute(breakerKey, async () => { + const controller = new AbortController(); + const timeoutMs = config.proxy?.timeoutMs || 30000; + const id = setTimeout(() => controller.abort(), timeoutMs); + + try { + const upstreamUrl = config.proxy?.upstreamUrl || 'http://localhost:4000'; + const url = `${upstreamUrl}/logs${req.params.endpoint ? `/${req.params.endpoint}` : ''}`; + + const response = await fetch(url, { + method: req.method, + headers: { + 'Content-Type': req.headers['content-type'] || 'application/json', + ...(req.headers['authorization'] ? { 'Authorization': req.headers['authorization'] as string } : {}), + }, + body: ['GET', 'HEAD'].includes(req.method) ? undefined : JSON.stringify(req.body), + signal: controller.signal, + }); + + if (!response.ok) { + throw new Error(`Upstream returned ${response.status}`); + } + + const contentType = response.headers.get('content-type'); + if (contentType && contentType.includes('application/json')) { + return await response.json(); + } else { + return await response.text(); + } + } finally { + clearTimeout(id); + } + }); + + if (typeof result === 'string') { + res.send(result); + } else { + res.json(result); + } + } catch (error) { + if (error && typeof error === 'object' && 'statusCode' in error && error.statusCode === 503) { + next(error); + } else { + next(new BadGatewayError(error instanceof Error ? error.message : 'Unknown error')); + } + } +}); + +export default logsRouter;