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
14 changes: 8 additions & 6 deletions backend/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,16 @@ API_BASE_URL=http://localhost:3001
# ─── Database ─────────────────────────────────────────────────────────────────
# Main PostgreSQL database connection URL
DATABASE_URL="postgresql://user:password@localhost:5433/flowfi?schema=public"
# Optional PostgreSQL read replica; omit to send all queries to the primary.
DATABASE_READ_REPLICA_URL=""

# PostgreSQL pool settings:
# Maximum database connections per backend process (default: 10)
PG_POOL_MAX=10
# How long an idle connection stays open before being closed (milliseconds, default: 30000)
PG_IDLE_TIMEOUT_MS=30000
# How long to wait when establishing a new database connection (milliseconds, default: 5000)
PG_CONNECTION_TIMEOUT_MS=5000
# Maximum database connections per backend process (default: 20)
PG_POOL_MAX=20
# How long an idle connection stays open before being closed (milliseconds, default: 10000)
PG_IDLE_TIMEOUT_MS=10000
# How long to wait when establishing a new database connection (milliseconds, default: 2000)
PG_CONNECTION_TIMEOUT_MS=2000
# Maximum time a PostgreSQL statement may run before cancellation (milliseconds, default: 30000)
PG_STATEMENT_TIMEOUT_MS=30000

Expand Down
11 changes: 4 additions & 7 deletions backend/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -28,9 +28,6 @@
"keywords": [],
"author": "",
"license": "ISC",
"overrides": {
"@types/pg": "8.20.0"
},
"dependencies": {
"@opentelemetry/api": "^1.9.1",
"@opentelemetry/auto-instrumentations-node": "^0.80.0",
Expand All @@ -39,8 +36,9 @@
"@opentelemetry/sdk-node": "^0.222.0",
"@opentelemetry/sdk-trace-node": "^2.11.0",
"@opentelemetry/semantic-conventions": "^1.43.0",
"@prisma/adapter-pg": "^6.19.3",
"@prisma/client": "^6.19.3",
"@prisma/adapter-pg": "^7.8.0",
"@prisma/client": "^7.8.0",
"@prisma/extension-read-replicas": "^0.5.0",
"@stellar/stellar-sdk": "^17.0.1",
"cors": "^2.8.6",
"dotenv": "^17.4.2",
Expand All @@ -59,14 +57,13 @@
"@types/eventsource": "^1.1.15",
"@types/express": "^5.0.6",
"@types/node": "^25.2.3",
"@types/pg": "8.20.0",
"@types/supertest": "^6.0.3",
"@types/swagger-jsdoc": "^6.0.4",
"@types/swagger-ui-express": "^4.1.6",
"@vitest/coverage-v8": "^3.2.7",
"eventsource": "^2.0.2",
"nodemon": "^3.1.11",
"prisma": "^6.19.3",
"prisma": "^7.4.1",
"supertest": "^7.1.0",
"ts-node": "^10.9.2",
"tsx": "^4.19.2",
Expand Down
1 change: 0 additions & 1 deletion backend/prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ generator client {

datasource db {
provider = "postgresql"
url = env("DATABASE_URL")
}

// User model - represents Stellar wallet addresses interacting with the protocol
Expand Down
26 changes: 13 additions & 13 deletions backend/src/controllers/stream.controller.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import type { Request, Response } from "express";
import { z } from "zod";
import { Prisma } from "../generated/prisma/index.js";
import { prisma } from "../lib/prisma.js";
import { prisma, withReplicaFallback } from "../lib/prisma.js";
import logger from "../logger.js";
import { claimableAmountService } from "../services/claimable.service.js";
import {
Expand Down Expand Up @@ -286,7 +286,7 @@ export const getStream = async (req: Request, res: Response) => {
return sendApiError(res, 400, "INVALID_STREAM_ID", "Invalid streamId parameter");
}

const stream = await prisma.stream.findUnique({
const stream = await withReplicaFallback((client) => client.stream.findUnique({
where: { streamId: parsedStreamId },
include: {
senderUser: true,
Expand All @@ -295,7 +295,7 @@ export const getStream = async (req: Request, res: Response) => {
orderBy: { timestamp: "desc" },
},
},
});
}));

if (!stream) {
// Fallback: try live RPC
Expand Down Expand Up @@ -385,8 +385,8 @@ export const getStreamEvents = async (req: Request, res: Response) => {
whereClause.eventType = eventType;
}

const [events, total] = await Promise.all([
prisma.streamEvent.findMany({
const [events, total] = await withReplicaFallback((client) => Promise.all([
client.streamEvent.findMany({
where: whereClause,
// `timestamp` is not unique (events in the same block/ledger can
// share a timestamp), so it can't be the sole sort key for cursor
Expand All @@ -396,8 +396,8 @@ export const getStreamEvents = async (req: Request, res: Response) => {
take: limit,
...(cursor ? { cursor: { id: cursor }, skip: 1 } : { skip: offset }),
}),
prisma.streamEvent.count({ where: whereClause }),
]);
client.streamEvent.count({ where: whereClause }),
]));

const hasMore = cursor
? events.length === limit
Expand Down Expand Up @@ -433,7 +433,7 @@ export const getStreamClaimableAmount = async (req: Request, res: Response) => {
}
}

const stream = await prisma.stream.findUnique({
const stream = await withReplicaFallback((client) => client.stream.findUnique({
where: { streamId: parsedStreamId },
select: {
streamId: true,
Expand All @@ -449,7 +449,7 @@ export const getStreamClaimableAmount = async (req: Request, res: Response) => {
totalPausedDuration: true,
updatedAt: true,
},
});
}));

if (!stream) {
// Fallback: try live RPC for claimable amount
Expand Down Expand Up @@ -522,8 +522,8 @@ export const getUserStreamSummary = async (
// unbounded DB queries. Power users with more than MAX_USER_STREAMS
// streams receive a truncated summary (the `truncated` flag lets the
// frontend offer a pagination/export fallback).
const [outgoingStreams, incomingStreams] = await Promise.all([
prisma.stream.findMany({
const [outgoingStreams, incomingStreams] = await withReplicaFallback((client) => Promise.all([
client.stream.findMany({
where: { sender: address },
orderBy: { startTime: "desc" },
take: MAX_USER_STREAMS,
Expand All @@ -542,7 +542,7 @@ export const getUserStreamSummary = async (
updatedAt: true,
},
}),
prisma.stream.findMany({
client.stream.findMany({
where: { recipient: address },
orderBy: { startTime: "desc" },
take: MAX_USER_STREAMS,
Expand All @@ -561,7 +561,7 @@ export const getUserStreamSummary = async (
updatedAt: true,
},
}),
]);
]));

const calculatedAt = Math.floor(nowMs / 1000);

Expand Down
6 changes: 3 additions & 3 deletions backend/src/lib/pg-pool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,9 @@ const parsePositiveIntegerEnv = (name: string, defaultValue: number): number =>

export const createPgPoolConfig = (overrides?: Partial<pg.PoolConfig>): pg.PoolConfig => ({
connectionString: process.env.DATABASE_URL,
max: parsePositiveIntegerEnv('PG_POOL_MAX', 10),
idleTimeoutMillis: parsePositiveIntegerEnv('PG_IDLE_TIMEOUT_MS', 30_000),
connectionTimeoutMillis: parsePositiveIntegerEnv('PG_CONNECTION_TIMEOUT_MS', 5_000),
max: parsePositiveIntegerEnv('PG_POOL_MAX', 20),
idleTimeoutMillis: parsePositiveIntegerEnv('PG_IDLE_TIMEOUT_MS', 10_000),
connectionTimeoutMillis: parsePositiveIntegerEnv('PG_CONNECTION_TIMEOUT_MS', 2_000),
statement_timeout: parsePositiveIntegerEnv('PG_STATEMENT_TIMEOUT_MS', 30_000),
...overrides,
});
Expand Down
75 changes: 68 additions & 7 deletions backend/src/lib/prisma.ts
Original file line number Diff line number Diff line change
@@ -1,29 +1,90 @@
import pg from 'pg';
import { PrismaPg } from '@prisma/adapter-pg';
import { readReplicas } from '@prisma/extension-read-replicas';
import { PrismaClient } from '../generated/prisma/index.js';
import { createPgPool } from './pg-pool.js';

const globalForPrisma = global as unknown as {
prisma?: PrismaClient;
pool?: pg.Pool;
readReplicaPool?: pg.Pool;
readReplicaClient?: PrismaClient;
primaryClient?: PrismaClient;
};

if (!globalForPrisma.pool) {
globalForPrisma.pool = createPgPool();
}

const adapter = new PrismaPg(globalForPrisma.pool);

export const prisma =
globalForPrisma.prisma ||
const log =
process.env.NODE_ENV === 'development'
? ['query', 'error', 'warn'] as const
: ['error'] as const;
const primaryClient =
globalForPrisma.primaryClient ||
new PrismaClient({
adapter,
log: process.env.NODE_ENV === 'development' ? ['query', 'error', 'warn'] : ['error'],
adapter: new PrismaPg(globalForPrisma.pool as unknown as ConstructorParameters<typeof PrismaPg>[0]),
log: [...log],
});
globalForPrisma.primaryClient = primaryClient;
const readReplicaUrl = process.env.DATABASE_READ_REPLICA_URL?.trim();

if (readReplicaUrl && !globalForPrisma.readReplicaClient) {
globalForPrisma.readReplicaPool = createPgPool({ connectionString: readReplicaUrl });
globalForPrisma.readReplicaClient = new PrismaClient({
adapter: new PrismaPg(globalForPrisma.readReplicaPool as unknown as ConstructorParameters<typeof PrismaPg>[0]),
log: [...log],
});
}

const client = readReplicaUrl && globalForPrisma.readReplicaClient
? primaryClient.$extends(readReplicas({ replicas: [globalForPrisma.readReplicaClient] }))
: primaryClient;

export const prisma = globalForPrisma.prisma || (client as PrismaClient);
if (process.env.NODE_ENV !== 'production') globalForPrisma.prisma = prisma;

const REPLICA_CONNECTION_ERROR_CODES = new Set([
'P1001', 'P1002', 'P1017', 'ECONNREFUSED', 'ECONNRESET', 'ETIMEDOUT',
'EHOSTUNREACH', 'ENETUNREACH', '57P01', '08000', '08001', '08003',
'08006', '08004', '08007', '08P01',
]);

function isReplicaConnectionError(error: unknown): boolean {
const pending: unknown[] = [error];
const seen = new Set<unknown>();
while (pending.length > 0) {
const current = pending.pop();
if (!current || typeof current !== 'object' || seen.has(current)) continue;
seen.add(current);
const details = current as { code?: unknown; cause?: unknown; message?: unknown; meta?: unknown; driverAdapterError?: unknown };
if (typeof details.code === 'string' && REPLICA_CONNECTION_ERROR_CODES.has(details.code)) return true;
if (typeof details.message === 'string' && /connection (?:terminated|closed|timeout)|connect(?:ion)? refused|server closed the connection|can't reach database server|failed to connect/i.test(details.message)) return true;
if (details.cause) pending.push(details.cause);
if (details.meta) pending.push(details.meta);
if (details.driverAdapterError) pending.push(details.driverAdapterError);
}
return false;
}

const REPLICA_FAILURE_COOLDOWN_MS = 5_000;
let replicaUnavailableUntil = 0;

/** Execute a read against the configured replica and retry connection failures on primary. */
export async function withReplicaFallback<T>(query: (client: PrismaClient) => Promise<T>): Promise<T> {
if (!readReplicaUrl) return query(prisma);
const primary = (prisma as PrismaClient & { $primary?: () => PrismaClient }).$primary?.();
if (!primary) return query(prisma);
if (Date.now() < replicaUnavailableUntil) return query(primary);
try {
return await query(prisma);
} catch (error) {
if (!isReplicaConnectionError(error)) throw error;
replicaUnavailableUntil = Date.now() + REPLICA_FAILURE_COOLDOWN_MS;
return query(primary);
}
}

export { getPoolMetrics } from './pg-pool.js';
export const pool = globalForPrisma.pool!;

export default prisma;
10 changes: 5 additions & 5 deletions backend/src/repositories/stream.repository.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { prisma } from '../lib/prisma.js';
import { prisma, withReplicaFallback } from '../lib/prisma.js';

/**
* Update the status and active flag of a stream in the database.
Expand Down Expand Up @@ -71,8 +71,8 @@ export const findStreams = async (params: FindStreamsParams): Promise<FindStream
const sortField = params.sortField || 'startTime';
const sortOrder = params.sortOrder || 'desc';

const [streams, total] = await Promise.all([
prisma.stream.findMany({
const [streams, total] = await withReplicaFallback((client) => Promise.all([
client.stream.findMany({
where,
orderBy: { [sortField]: sortOrder },
take: params.limit,
Expand All @@ -82,8 +82,8 @@ export const findStreams = async (params: FindStreamsParams): Promise<FindStream
recipientUser: true,
},
}),
prisma.stream.count({ where }),
]);
client.stream.count({ where }),
]));

return { streams, total, hasMore: params.offset + streams.length < total };
};
10 changes: 5 additions & 5 deletions backend/src/repositories/streamEvent.repository.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { prisma } from "../lib/prisma.js";
import { withReplicaFallback } from "../lib/prisma.js";
import type { Prisma } from "../generated/prisma/index.js";

/**
Expand Down Expand Up @@ -136,16 +136,16 @@ export async function listEventsForWallet(
where.eventType = { in: types };
}

const [events, total] = await Promise.all([
prisma.streamEvent.findMany({
const [events, total] = await withReplicaFallback((client) => Promise.all([
client.streamEvent.findMany({
where,
orderBy: { timestamp: "desc" },
skip: offset,
take: limit,
...(includeStream ? { include: { stream: true } } : {}),
}),
prisma.streamEvent.count({ where }),
]);
client.streamEvent.count({ where }),
]));

return {
events,
Expand Down
2 changes: 2 additions & 0 deletions backend/tests/auth.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,12 @@ import * as StellarSdk from '@stellar/stellar-sdk';
import express from 'express';
import { rateLimit } from 'express-rate-limit';
import app from '../src/app.js';
import { prisma } from '../src/lib/prisma.js';
import { __authChallengeTestUtils, requireAdmin, signJwt } from '../src/middleware/auth.js';

// Mocking prisma for any downstream dependency
vi.mock('../src/lib/prisma.js', () => ({
withReplicaFallback: (query: (client: any) => unknown) => Promise.resolve(query(prisma)),
default: {
stream: { findMany: vi.fn(() => Promise.resolve([])) },
streamEvent: {
Expand Down
1 change: 1 addition & 0 deletions backend/tests/integration/events-list.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ const mocks = vi.hoisted(() => ({
}));

vi.mock('../../src/lib/prisma.js', () => ({
withReplicaFallback: (query: (client: any) => unknown) => Promise.resolve(query(mocks.prisma)),
default: mocks.prisma,
prisma: mocks.prisma,
}));
Expand Down
1 change: 1 addition & 0 deletions backend/tests/integration/indexer-worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ const { mockPrisma, mockSseService } = vi.hoisted(() => ({
}));

vi.mock('../../src/lib/prisma.js', () => ({
withReplicaFallback: (query: (client: any) => unknown) => Promise.resolve(query(mockPrisma)),
prisma: mockPrisma,
default: mockPrisma,
}));
Expand Down
1 change: 1 addition & 0 deletions backend/tests/integration/streams.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ vi.mock('../../src/lib/redis.js', () => ({
}));

vi.mock('../../src/lib/prisma.js', () => ({
withReplicaFallback: (query: (client: any) => unknown) => Promise.resolve(query(mockPrisma)),
prisma: mockPrisma,
}));

Expand Down
Loading
Loading