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
58 changes: 58 additions & 0 deletions backend/src/lib/indexer-state.ts
Original file line number Diff line number Diff line change
@@ -1 +1,59 @@
import { prisma } from './prisma.js';

Check failure on line 1 in backend/src/lib/indexer-state.ts

View workflow job for this annotation

GitHub Actions / Backend npm test

tests/indexer-state.test.ts

Error: [vitest] There was an error when mocking a module. If you are using "vi.mock" factory, make sure there are no top level variables inside, since this call is hoisted to top of the file. Read more: https://vitest.dev/api/vi.html#vi-mock ❯ src/lib/indexer-state.ts:1:1 Caused by: Caused by: ReferenceError: Cannot access 'mockFindUnique' before initialization ❯ tests/indexer-state.test.ts:9:19 ❯ src/lib/indexer-state.ts:1:1
import logger from '../logger.js';

export const INDEXER_STATE_ID = 'singleton';

export interface IndexerStateRow {
id: string;
lastLedger: number;
lastCursor: string | null;
createdAt: Date;
updatedAt: Date;
}

/**
* Ensure the singleton indexer_state row exists.
* Uses a catch-and-retry pattern to handle the race condition where two
* concurrent callers attempt the first insert simultaneously. If the
* unique-constraint violation fires, we treat it as success and re-read.
*/
export async function ensureIndexerState(
startLedger: number,
): Promise<IndexerStateRow> {
const existing = await prisma.indexerState.findUnique({
where: { id: INDEXER_STATE_ID },
});
if (existing) return existing;

try {
const created = await prisma.indexerState.create({
data: {
id: INDEXER_STATE_ID,
lastLedger: startLedger,
lastCursor: null,
},
});
return created;
} catch (err: unknown) {
// P2002 = Prisma unique-constraint violation (code "P2002")
if (
err instanceof Error &&
'code' in err &&
(err as { code: string }).code === 'P2002'
) {
logger.warn(
'[IndexerState] Concurrent first-insert detected; re-reading existing row.',
);
const existingAfterRace = await prisma.indexerState.findUnique({
where: { id: INDEXER_STATE_ID },
});
if (!existingAfterRace) {
throw new Error(
'[IndexerState] Unique-constraint violation but row not found after race.',
);
}
return existingAfterRace;
}
throw err;
}
}
14 changes: 14 additions & 0 deletions backend/src/lib/pg-pool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,3 +19,17 @@ export const createPgPoolConfig = (): pg.PoolConfig => ({
});

export const createPgPool = () => new pg.Pool(createPgPoolConfig());

export interface PoolMetrics {
totalCount: number;
idleCount: number;
waitingCount: number;
}

export function getPoolMetrics(pool: pg.Pool): PoolMetrics {
return {
totalCount: pool.totalCount,
idleCount: pool.idleCount,
waitingCount: pool.waitingCount,
};
}
3 changes: 3 additions & 0 deletions backend/src/lib/prisma.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,4 +23,7 @@ export const prisma =

if (process.env.NODE_ENV !== 'production') globalForPrisma.prisma = prisma;

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

export default prisma;
29 changes: 29 additions & 0 deletions backend/src/middleware/admin-rate-limiter.middleware.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
import { rateLimit } from 'express-rate-limit';
import type { Request, Response } from 'express';

export const adminRateLimiter = rateLimit({
windowMs: 1 * 60 * 1000, // 1 minute
max: 30, // Stricter limit: 30 requests per minute for admin endpoints
standardHeaders: true,
legacyHeaders: false,
message: {
error: 'Too many admin requests',
message: 'You have exceeded the admin rate limit. Please try again later.',
status: 429,
},
keyGenerator: (req: Request): string => {
// Use x-forwarded-for or remote address as key
const forwarded = req.headers['x-forwarded-for'];
if (typeof forwarded === 'string') {
return forwarded.split(',')[0].trim();
}
return req.ip ?? 'unknown';
},
skip: (req: Request): boolean => {
// Skip rate limiting in test environment
return process.env.NODE_ENV === 'test';
},
handler: (req: Request, res: Response, _next, options): void => {
res.status(options.statusCode).json(options.message);
},
});
6 changes: 5 additions & 1 deletion backend/src/routes/v1/admin.routes.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,15 @@
import { Router } from 'express';
import type { Request, Response } from 'express';
import { requireAdmin } from '../../middleware/auth.js';
import { adminRateLimiter } from '../../middleware/admin-rate-limiter.middleware.js';
import {
getIndexerStatus,
resetIndexer,
replayFromLedger,
} from '../../services/indexerService.js';

import { prisma } from '../../lib/prisma.js';
import { prisma, pool } from '../../lib/prisma.js';
import { getPoolMetrics } from '../../lib/pg-pool.js';
import { INDEXER_STATE_ID } from '../../lib/indexer-state.js';
import { sseService } from '../../services/sse.service.js';
import { cache } from '../../lib/redis.js';
Expand All @@ -17,6 +19,7 @@ const router = Router();

// All admin routes require admin JWT
router.use(requireAdmin);
router.use(adminRateLimiter);

/**
* @openapi
Expand Down Expand Up @@ -131,6 +134,7 @@ async function buildAdminMetrics() {
},
sse: { activeConnections: sseService.getClientCount() },
cache: cache.getStats(),
pgPool: getPoolMetrics(pool),
indexer: {
lastLedger: indexerState?.lastLedger ?? 0,
lagSeconds,
Expand Down
12 changes: 2 additions & 10 deletions backend/src/workers/soroban-event-worker.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { rpc, xdr, StrKey } from "@stellar/stellar-sdk";
import { prisma } from "../lib/prisma.js";
import { INDEXER_STATE_ID } from "../lib/indexer-state.js";
import { INDEXER_STATE_ID, ensureIndexerState } from "../lib/indexer-state.js";
import { sseService } from "../services/sse.service.js";
import logger from "../logger.js";
import { Prisma } from "../generated/prisma/index.js";
Expand Down Expand Up @@ -33,7 +33,7 @@
* Full value = hi * 2^64 + lo.
*/
export function decodeI128(val: xdr.ScVal): string {
const parts = val.i128();

Check failure on line 36 in backend/src/workers/soroban-event-worker.ts

View workflow job for this annotation

GitHub Actions / Backend npm test

tests/integration/stream-lifecycle.test.ts > Stream Lifecycle Integration Tests > Full lifecycle: create → top up → partial withdraw → cancel > walks a single stream through every phase and verifies indexer state

TypeError: i128 not set ❯ ChildUnion.get ../node_modules/@stellar/js-xdr/lib/webpack:/XDR/src/union.js:27:13 ❯ ChildUnion.get [as i128] ../node_modules/@stellar/js-xdr/lib/webpack:/XDR/src/union.js:171:25 ❯ decodeI128 src/workers/soroban-event-worker.ts:36:21 ❯ SorobanEventWorker.handleStreamCreated src/workers/soroban-event-worker.ts:459:27 ❯ SorobanEventWorker.processEvent src/workers/soroban-event-worker.ts:297:20 ❯ tests/integration/stream-lifecycle.test.ts:635:20
const hi = BigInt.asIntN(64, BigInt(parts.hi().toString()));
const lo = BigInt.asUintN(64, BigInt(parts.lo().toString()));
return ((hi << 64n) | lo).toString();
Expand Down Expand Up @@ -189,15 +189,7 @@
*/
private async fetchAndProcessEvents(): Promise<void> {
// Ensure an IndexerState row exists on first run.
const state = await prisma.indexerState.upsert({
where: { id: INDEXER_STATE_ID },
create: {
id: INDEXER_STATE_ID,
lastLedger: this.startLedger,
lastCursor: null,
},
update: {},
});
const state = await ensureIndexerState(this.startLedger);

const baseFilter = {
filters: [
Expand Down
21 changes: 21 additions & 0 deletions backend/tests/admin-rate-limiter.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import { describe, it, expect, vi, beforeEach } from 'vitest';

vi.mock('../src/lib/prisma.js', () => ({
prisma: {},
}));

vi.mock('../src/logger.js', () => ({
default: { info: vi.fn(), error: vi.fn(), warn: vi.fn() },
}));

import { adminRateLimiter } from '../src/middleware/admin-rate-limiter.middleware.js';

describe('adminRateLimiter', () => {
beforeEach(() => {
vi.clearAllMocks();
});

it('is a function', () => {
expect(typeof adminRateLimiter).toBe('function');
});
});
98 changes: 98 additions & 0 deletions backend/tests/indexer-state.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
import { describe, it, expect, vi, beforeEach } from 'vitest';

const mockFindUnique = vi.fn();
const mockCreate = vi.fn();

vi.mock('../src/lib/prisma.js', () => ({
prisma: {
indexerState: {
findUnique: mockFindUnique,
create: mockCreate,
},
},
}));

vi.mock('../src/logger.js', () => ({
default: {
warn: vi.fn(),
info: vi.fn(),
error: vi.fn(),
},
}));

import { ensureIndexerState } from '../src/lib/indexer-state.js';

describe('ensureIndexerState', () => {
beforeEach(() => {
vi.clearAllMocks();
});

it('returns existing row if it already exists', async () => {
const existing = {
id: 'singleton',
lastLedger: 42,
lastCursor: 'cur',
createdAt: new Date(),
updatedAt: new Date(),
};
mockFindUnique.mockResolvedValueOnce(existing);

const result = await ensureIndexerState(0);

expect(result).toEqual(existing);
expect(mockCreate).not.toHaveBeenCalled();
});

it('creates and returns a new row when none exists', async () => {
const created = {
id: 'singleton',
lastLedger: 10,
lastCursor: null,
createdAt: new Date(),
updatedAt: new Date(),
};
mockFindUnique.mockResolvedValueOnce(null);
mockCreate.mockResolvedValueOnce(created);

const result = await ensureIndexerState(10);

expect(result).toEqual(created);
expect(mockCreate).toHaveBeenCalledWith({
data: { id: 'singleton', lastLedger: 10, lastCursor: null },
});
});

it('re-reads the row on unique-constraint violation (P2002) without throwing', async () => {
const existingRow = {
id: 'singleton',
lastLedger: 5,
lastCursor: null,
createdAt: new Date(),
updatedAt: new Date(),
};

// First findUnique returns null (no row yet)
mockFindUnique.mockResolvedValueOnce(null);
// Create throws P2002 (race condition duplicate insert)
const p2002Error = Object.assign(new Error('Unique constraint failed'), {
code: 'P2002',
});
mockCreate.mockRejectedValueOnce(p2002Error);
// Second findUnique returns the existing row created by the concurrent caller
mockFindUnique.mockResolvedValueOnce(existingRow);

const result = await ensureIndexerState(0);

expect(result).toEqual(existingRow);
expect(mockCreate).toHaveBeenCalledTimes(1);
expect(mockFindUnique).toHaveBeenCalledTimes(2);
});

it('re-throws non-P2002 errors', async () => {
mockFindUnique.mockResolvedValueOnce(null);
const genericError = new Error('connection refused');
mockCreate.mockRejectedValueOnce(genericError);

await expect(ensureIndexerState(0)).rejects.toThrow('connection refused');
});
});
18 changes: 18 additions & 0 deletions backend/tests/pg-pool.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,4 +45,22 @@ describe('pg-pool', () => {
}),
);
});

it('getPoolMetrics returns totalCount, idleCount, and waitingCount from the pool', async () => {
const { getPoolMetrics } = await import('../src/lib/pg-pool.js');

const mockPool = {
totalCount: 10,
idleCount: 5,
waitingCount: 2,
} as any;

const metrics = getPoolMetrics(mockPool);

expect(metrics).toEqual({
totalCount: 10,
idleCount: 5,
waitingCount: 2,
});
});
});
Loading