Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 52 additions & 0 deletions docs/api-logs.md
Original file line number Diff line number Diff line change
@@ -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"
}
}
```
36 changes: 36 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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
Expand Down
13 changes: 5 additions & 8 deletions src/lib/circuitBreaker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ export interface CircuitBreakerConfig {
failureThreshold?: number;
cooldownMs?: number;
successThreshold?: number;
onOpenError?: (message: string) => Error;
}

export interface CircuitBreakerMetrics {
Expand Down Expand Up @@ -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);
}
Expand Down
58 changes: 58 additions & 0 deletions src/routes/logs.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
68 changes: 68 additions & 0 deletions src/routes/logs.ts
Original file line number Diff line number Diff line change
@@ -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;
Loading