Skip to content
Closed
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
6 changes: 6 additions & 0 deletions ENV_VARS.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,10 @@
| VAPID_PRIVATE_KEY | VAPID private key for Web Push API | auto-generated | No |
| WS_ENABLED | Enable/disable WebSocket support | true | No |
| WS_PORT | WebSocket port | 3001 | No |
| DB_READ_REPLICA_URLS | Comma-separated PostgreSQL read replica URLs | - | No |
| DB_REPLICA_MAX_LAG_MS | Maximum replica lag before primary failover | 5000 | No |
| DB_REPLICA_HEALTH_CHECK_INTERVAL_MS | Interval for replica health checks | 30000 | No |
| DB_REPLICA_FAILOVER_COOLDOWN_MS | Cooldown after replica failover | 15000 | No |

## Frontend

Expand All @@ -37,6 +41,8 @@ AGENTICPAY_ALLOWED_SIGNATURE_ORIGINS=https://agenticpay.com,http://localhost:300
VAPID_PUBLIC_KEY=your-vapid-public-key
VAPID_PRIVATE_KEY=your-vapid-private-key
WS_ENABLED=true
DB_READ_REPLICA_URLS=
DB_REPLICA_MAX_LAG_MS=5000
```

- `.env.development` — local development
Expand Down
89 changes: 75 additions & 14 deletions PERFORMANCE_GUIDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -65,22 +65,20 @@ Reduces payload size by 60-80% on average, improving:

### Compression Methods

1. **Brotli** (preferred): 20-30% smaller than gzip
- Quality level: 5 (balanced speed/compression)
- Mode: Text optimization
AgenticPay uses the maintained Express `compression` middleware with streaming
backpressure support, negotiated encodings, and route-level filters.

2. **Gzip** (fallback): Universal support
- Compression level: 6
- Minimum size threshold: 1KB
- Compression level: 6
- Minimum size threshold: 1KB by default
- Skips images, audio, video, archives, and already-compressed content

### Implementation

Located in: `backend/src/middleware/compression.ts`

```typescript
app.use(compressionMiddleware({
brotliLevel: 5,
gzipLevel: 6,
level: 6,
minSizeBytes: 1024,
}));
```
Expand All @@ -90,7 +88,7 @@ app.use(compressionMiddleware({
Access compression metrics via:

```
GET /api/v1/monitoring/pool/compression
GET /api/v1/monitoring/compression
```

Returns:
Expand All @@ -108,6 +106,35 @@ Returns:

---

## API Response Streaming

Large exports are streamed with chunked transfer encoding instead of being buffered in memory.

Endpoints:

```bash
GET /api/v1/exports/audit/stream?format=csv&limit=100000
GET /api/v1/exports/audit/stream?format=jsonl&batchSize=1000
GET /api/v1/exports/payments/stream?format=csv
```

Reusable helpers live in `backend/src/middleware/streaming.ts`:

```typescript
const query = parseStreamingQuery(req.query);
await streamDataset({
req,
res,
items: takeStreamItems(fetchRows(), query.limit),
format: query.format,
});
```

The streaming helpers set `Transfer-Encoding: chunked`, disable proxy buffering with
`X-Accel-Buffering: no`, honor HTTP backpressure, and track completed, aborted, and failed streams.

---

## Cursor-Based Pagination

### Overview
Expand Down Expand Up @@ -198,6 +225,40 @@ GET /api/v1/payments -H "If-None-Match: abc123def456"

Optimized connection pooling with PgBouncer for efficient resource utilization:

### Read Replicas and Failover

Read replica routing is configured with:

```bash
DB_READ_REPLICA_URLS=postgresql://user:pass@replica-a:5432/agenticpay,postgresql://user:pass@replica-b:5432/agenticpay
DB_REPLICA_MAX_LAG_MS=5000
DB_REPLICA_HEALTH_CHECK_INTERVAL_MS=30000
DB_REPLICA_FAILOVER_COOLDOWN_MS=15000
```

`backend/src/config/database.ts` exposes `ReadReplicaRouter`, which routes `SELECT` and `WITH`
queries across healthy replicas and falls back to `DATABASE_URL` when no replica is available or
replica lag exceeds the configured threshold. Terraform can provision replicas with
`db_read_replica_count` and wires `DB_READ_REPLICA_URLS` into the backend service.

### WebSocket Pooling

WebSocket connections are managed by `backend/src/websocket/pool.ts`. The pool enforces capacity,
tracks active and queued connections, batches outbound messages through `ManagedConnection`, and
supports clean shutdown. Tune batching with:

```typescript
attachWebSocketServer({
server,
options: {
maxConnections: 250,
maxQueueSizePerConnection: 500,
flushIntervalMs: 25,
maxBatchSize: 50,
},
});
```

**Benefits:**
- Prevents connection exhaustion
- Detects and prevents connection leaks
Expand Down Expand Up @@ -239,7 +300,7 @@ Located in: `backend/src/config/database.ts`
Access pool health via:

```
GET /api/v1/monitoring/pool/health
GET /api/v1/monitoring/health
```

Returns:
Expand All @@ -261,7 +322,7 @@ Returns:
Automatic detection of connection leaks:

```
GET /api/v1/monitoring/pool/leaks
GET /api/v1/monitoring/leaks
```

- Monitors connection acquisition/release
Expand All @@ -271,7 +332,7 @@ GET /api/v1/monitoring/pool/leaks
### Metrics Endpoint

```
GET /api/v1/monitoring/pool/metrics
GET /api/v1/monitoring/metrics
```

Returns comprehensive pool statistics including:
Expand Down Expand Up @@ -358,7 +419,7 @@ cache.registerWarmer('dashboard:overview',
### Metrics Endpoint

```
GET /api/v1/monitoring/pool/cache
GET /api/v1/monitoring/cache
```

Returns:
Expand All @@ -383,7 +444,7 @@ Returns:
Comprehensive view of all performance metrics:

```
GET /api/v1/monitoring/pool/performance
GET /api/v1/monitoring/performance
```

Returns combined metrics:
Expand Down
6 changes: 6 additions & 0 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,12 @@ RATE_LIMIT_ENTERPRISE=1000
RATE_LIMIT_WINDOW_MS=900000
COMPRESSION_THRESHOLD=1024

# Database read replicas (comma-separated PostgreSQL URLs)
DB_READ_REPLICA_URLS=
DB_REPLICA_MAX_LAG_MS=5000
DB_REPLICA_HEALTH_CHECK_INTERVAL_MS=30000
DB_REPLICA_FAILOVER_COOLDOWN_MS=15000

# Security headers
HSTS_MAX_AGE_SECONDS=31536000
PERMISSIONS_POLICY=camera=(), microphone=(), geolocation=(), payment=(), usb=(), magnetometer=(), gyroscope=(), interest-cohort=()
Expand Down
10 changes: 10 additions & 0 deletions backend/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,10 @@ const envSchema = z.object({
DB_POOL_ACQUIRE_TIMEOUT_MS: z.string().default('30000'),
DB_POOL_MAX_USES: z.string().default('7500'),
DB_STATEMENT_TIMEOUT_MS: z.string().default('30000'),
DB_READ_REPLICA_URLS: z.string().default(''),
DB_REPLICA_MAX_LAG_MS: z.string().default('5000'),
DB_REPLICA_HEALTH_CHECK_INTERVAL_MS: z.string().default('30000'),
DB_REPLICA_FAILOVER_COOLDOWN_MS: z.string().default('15000'),
HSTS_MAX_AGE_SECONDS: z.string().default('31536000'),
PERMISSIONS_POLICY: z
.string()
Expand Down Expand Up @@ -83,6 +87,12 @@ export const config = {
maxUses: Number(env.DB_POOL_MAX_USES),
statementTimeoutMs: Number(env.DB_STATEMENT_TIMEOUT_MS),
},
replicas: {
urls: env.DB_READ_REPLICA_URLS.split(',').map((url) => url.trim()).filter(Boolean),
maxLagMs: Number(env.DB_REPLICA_MAX_LAG_MS),
healthCheckIntervalMs: Number(env.DB_REPLICA_HEALTH_CHECK_INTERVAL_MS),
failoverCooldownMs: Number(env.DB_REPLICA_FAILOVER_COOLDOWN_MS),
},
},
security: {
hsts: {
Expand Down
70 changes: 70 additions & 0 deletions backend/src/config/database-read-replica.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
import { describe, expect, it } from 'vitest';
import {
ReadReplicaRouter,
buildReplicaConfigs,
isReadQuery,
} from './database';

describe('read replica routing', () => {
it('detects read queries conservatively', () => {
expect(isReadQuery('select * from payments')).toBe(true);
expect(isReadQuery(' WITH recent AS (select 1) select * from recent')).toBe(true);
expect(isReadQuery('update payments set status = $1')).toBe(false);
});

it('routes writes to primary and reads to healthy replicas', () => {
const router = new ReadReplicaRouter(
['postgres://replica-1/db', 'postgres://replica-2/db'],
'postgres://primary/db',
5000,
);

expect(router.select('UPDATE payments SET status = $1')).toEqual({
url: 'postgres://primary/db',
source: 'primary',
reason: 'write_query',
});

expect(router.select('SELECT * FROM payments')).toEqual({
url: 'postgres://replica-1/db',
source: 'replica',
reason: 'healthy_replica',
});
expect(router.select('SELECT * FROM invoices')).toMatchObject({
url: 'postgres://replica-2/db',
source: 'replica',
});
});

it('fails over to primary when all replicas are unhealthy or lagging', () => {
const router = new ReadReplicaRouter(
['postgres://replica-1/db', 'postgres://replica-2/db'],
'postgres://primary/db',
100,
);

router.updateHealth('postgres://replica-1/db', { healthy: false });
router.updateHealth('postgres://replica-2/db', { healthy: true, lagMs: 500 });

expect(router.select('SELECT * FROM payments')).toEqual({
url: 'postgres://primary/db',
source: 'primary',
reason: 'replica_unavailable',
});
});

it('builds replica configs from environment URLs', () => {
const previous = process.env.DB_READ_REPLICA_URLS;
process.env.DB_READ_REPLICA_URLS = 'postgres://user:pass@replica-a:5432/app, postgres://user:pass@replica-b/app';

try {
expect(buildReplicaConfigs()).toMatchObject([
{ host: 'replica-a', port: 5432, database: 'app', user: 'user', enabled: true },
{ host: 'replica-b', port: 5432, database: 'app', user: 'user', enabled: true },
]);
} finally {
if (previous === undefined) delete process.env.DB_READ_REPLICA_URLS;
else process.env.DB_READ_REPLICA_URLS = previous;
}
});
});
77 changes: 77 additions & 0 deletions backend/src/config/database.ts
Original file line number Diff line number Diff line change
Expand Up @@ -874,6 +874,83 @@ export function isReadQuery(sql: string): boolean {
return /^\s*(SELECT|WITH\s)/i.test(sql);
}

export type ReplicaHealth = "healthy" | "lagging" | "unhealthy";

export interface ReadReplicaTarget {
url: string;
health: ReplicaHealth;
lagMs: number;
lastCheckedAt: number;
failureCount: number;
}

export interface ReplicaSelection {
url: string;
source: "primary" | "replica";
reason: "write_query" | "no_replicas" | "healthy_replica" | "replica_unavailable";
}

export class ReadReplicaRouter {
private replicas: ReadReplicaTarget[];
private nextReplicaIndex = 0;

constructor(
replicaUrls = buildReplicaUrls(),
private readonly primaryUrl = process.env.DATABASE_URL ?? "",
private readonly maxLagMs = envInt("DB_REPLICA_MAX_LAG_MS", 5000),
) {
this.replicas = replicaUrls.map((url) => ({
url,
health: "healthy",
lagMs: 0,
lastCheckedAt: 0,
failureCount: 0,
}));
}

select(sql: string): ReplicaSelection {
if (!isReadQuery(sql)) {
return { url: this.primaryUrl, source: "primary", reason: "write_query" };
}

const healthyReplicas = this.replicas.filter(
(replica) => replica.health === "healthy" && replica.lagMs <= this.maxLagMs,
);

if (this.replicas.length === 0) {
return { url: this.primaryUrl, source: "primary", reason: "no_replicas" };
}

if (healthyReplicas.length === 0) {
return { url: this.primaryUrl, source: "primary", reason: "replica_unavailable" };
}

const replica = healthyReplicas[this.nextReplicaIndex % healthyReplicas.length];
this.nextReplicaIndex = (this.nextReplicaIndex + 1) % healthyReplicas.length;
return { url: replica.url, source: "replica", reason: "healthy_replica" };
}

updateHealth(url: string, params: { healthy: boolean; lagMs?: number; checkedAt?: number }): void {
const replica = this.replicas.find((candidate) => candidate.url === url);
if (!replica) return;

replica.lagMs = params.lagMs ?? replica.lagMs;
replica.lastCheckedAt = params.checkedAt ?? Date.now();
replica.health = !params.healthy
? "unhealthy"
: replica.lagMs > this.maxLagMs
? "lagging"
: "healthy";
replica.failureCount = replica.health === "healthy" ? 0 : replica.failureCount + 1;
}

snapshot(): ReadReplicaTarget[] {
return this.replicas.map((replica) => ({ ...replica }));
}
}

export const readReplicaRouter = new ReadReplicaRouter();

// ── Query Profiler ────────────────────────────────────────────────────────────

export interface QueryProfile {
Expand Down
Loading
Loading