Skip to content
Open
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
2 changes: 2 additions & 0 deletions backend/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@ import v1Routes from "./routes/v1/index.js";
import healthRoutes from "./routes/health.routes.js";
import "./lib/stream-id.js";



const app = express();
const isProduction = process.env.NODE_ENV === "production";
const rawCors = process.env.CORS_ALLOWED_ORIGINS ?? "";
Expand Down
1 change: 0 additions & 1 deletion backend/src/lib/indexer-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ export interface IndexerStateRow {
id: string;
lastLedger: number;
lastCursor: string | null;
createdAt: Date;
updatedAt: Date;
}

Expand Down
10 changes: 5 additions & 5 deletions backend/src/middleware/admin-rate-limiter.middleware.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { rateLimit } from 'express-rate-limit';
import { rateLimit, ipKeyGenerator } from 'express-rate-limit';
import type { Request, Response } from 'express';

export const adminRateLimiter = rateLimit({
Expand All @@ -11,13 +11,13 @@ export const adminRateLimiter = rateLimit({
message: 'You have exceeded the admin rate limit. Please try again later.',
status: 429,
},
keyGenerator: (req: Request): string => {
keyGenerator: (req: Request, res: Response): 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();
if (typeof forwarded === 'string' && forwarded.trim()) {
return forwarded.split(',')[0]?.trim() || ipKeyGenerator(req.ip ?? 'unknown');
}
return req.ip ?? 'unknown';
return ipKeyGenerator(req.ip ?? 'unknown');
},
skip: (req: Request): boolean => {
// Skip rate limiting in test environment
Expand Down
20 changes: 16 additions & 4 deletions backend/src/middleware/rate-limiter.middleware.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,22 @@
import { rateLimit } from 'express-rate-limit';
import { rateLimit, type Options } from 'express-rate-limit';

export const globalRateLimiter = rateLimit({
/**
* Shared factory to create an express-rate-limit instance with common configuration.
*
* @param options Configuration options for express-rate-limit
* @returns Express rate limit middleware
*/
export function createRateLimiter(options: Partial<Options>) {
return rateLimit({
standardHeaders: true, // Return rate limit info in the `RateLimit-*` headers
legacyHeaders: false, // Disable the `X-RateLimit-*` headers
...options,
});
}

export const globalRateLimiter = createRateLimiter({
windowMs: 1 * 60 * 1000, // 1 minute
max: 100, // Limit each IP to 100 requests per `window` (here, per minute)
standardHeaders: true, // Return rate limit info in the `RateLimit-*` headers
legacyHeaders: false, // Disable the `X-RateLimit-*` headers
message: {
message: 'Too many requests, please try again later.',
status: 429,
Expand Down
6 changes: 2 additions & 4 deletions backend/src/middleware/stream-rate-limiter.middleware.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { rateLimit } from 'express-rate-limit';
import { createRateLimiter } from './rate-limiter.middleware.js';
import { type Request, type Response, type NextFunction } from 'express';
import type { AuthenticatedRequest } from '../types/auth.types.js';
import logger from '../logger.js';
Expand All @@ -22,11 +22,9 @@ export function createStreamRateLimiter(
// Read from environment variable, default to 10 if not set
const max = options?.max ?? (process.env.STREAM_CREATE_RATE_LIMIT ? parseInt(process.env.STREAM_CREATE_RATE_LIMIT, 10) : 10);

return rateLimit({
return createRateLimiter({
windowMs,
max,
standardHeaders: true, // Return rate limit info in the `RateLimit-*` headers
legacyHeaders: false, // Disable the `X-RateLimit-*` headers
message: {
error: 'Too many stream creation requests - rate limit exceeded',
message: 'You have exceeded the rate limit for stream creation. Please try again later.',
Expand Down
4 changes: 3 additions & 1 deletion backend/src/routes/health.routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ import { Router, type Request, type Response } from 'express';
import { prisma } from '../lib/prisma.js';
import { INDEXER_STATE_ID } from '../lib/indexer-state.js';
import { sorobanEventWorker } from '../workers/soroban-event-worker.js';
import { isRedisAvailable } from '../lib/redis.js';
import { checkRpcHealth } from '../services/sorobanService.js';

const router = Router();

Expand Down Expand Up @@ -161,7 +163,7 @@ router.get('/', async (_req: Request, res: Response) => {
status: dbStatus === 'connected' ? 'ok' : 'down',
},
indexer: {
status: !indexerEnabled ? 'disabled' : indexerDegraded ? 'degraded' : 'ok',
status: !indexerEnabled ? 'disabled' : (indexerLagDegraded || indexerFailureDegraded) ? 'degraded' : 'ok',
enabled: indexerEnabled,
lagSeconds: indexerLag === -1 ? null : indexerLag,
},
Expand Down
3 changes: 1 addition & 2 deletions backend/src/services/soroban-indexer.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -185,8 +185,7 @@ export class SorobanIndexerService {
const ratePerSecond = this.readString(value, 'rate_per_second', 'ratePerSecond');
const depositedAmount = this.readString(value, 'deposited_amount', 'depositedAmount');
const startTimeRaw = value.start_time ?? value.startTime ?? timestamp;
const startTime = BigInt(startTimeRaw ?? timestamp);
const timestampBigInt = BigInt(timestamp);
const startTime = BigInt(startTimeRaw as string | number);

if (!sender || !recipient || !tokenAddress || !ratePerSecond || !depositedAmount) return;

Expand Down
18 changes: 6 additions & 12 deletions backend/src/workers/soroban-event-worker.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,8 @@
import { randomUUID } from "crypto";
import { rpc, xdr, StrKey } from "@stellar/stellar-sdk";
import { prisma } from "../lib/prisma.js";
import { INDEXER_STATE_ID, ensureIndexerState } from "../lib/indexer-state.js";
import { sseService } from "../services/sse.service.js";
import logger, { requestContext } from "../logger.js";
import logger from "../logger.js";
import { Prisma } from "../generated/prisma/index.js";
import "../lib/stream-id.js";

Expand Down Expand Up @@ -113,12 +112,7 @@ export class SorobanEventWorker {
/** Recent attempt outcomes for sliding-window spike detection. */
private recentOutcomes: { ok: boolean; at: number }[] = [];

/**
* Stable id attached to every log line emitted by the background poll
* loop, since these callbacks fire outside of any HTTP request and would
* otherwise have no requestContext (and thus no correlation id) at all.
*/
private readonly workerId = `soroban-worker:${randomUUID()}`;


constructor() {
const rpcUrl =
Expand Down Expand Up @@ -1126,13 +1120,13 @@ export class SorobanEventWorker {
});

// Calculate the duration of this pause interval
let additionalPausedDuration = 0;
let additionalPausedDuration = 0n;
if (currentStream.pausedAt) {
additionalPausedDuration = timestamp - currentStream.pausedAt;
additionalPausedDuration = BigInt(timestamp) - BigInt(currentStream.pausedAt);
}

const newTotalPausedDuration =
currentStream.totalPausedDuration + additionalPausedDuration;
Number(BigInt(currentStream.totalPausedDuration) + additionalPausedDuration);

await tx.stream.update({
where: { streamId },
Expand Down Expand Up @@ -1175,7 +1169,7 @@ export class SorobanEventWorker {
metadata: JSON.stringify({
sender,
newEndTime,
pausedDuration: additionalPausedDuration,
pausedDuration: Number(additionalPausedDuration),
totalPausedDuration: newTotalPausedDuration,
}),
},
Expand Down
32 changes: 16 additions & 16 deletions backend/tests/claimable.service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,12 @@ import { ClaimableAmountService } from '../src/services/claimable.service.js';

function makeStreamState(overrides: Partial<Parameters<ClaimableAmountService['getClaimableAmount']>[0]> = {}) {
return {
streamId: 1,
streamId: 1n,
ratePerSecond: '10',
depositedAmount: '100',
withdrawnAmount: '0',
lastUpdateTime: 0,
startTime: 0,
lastUpdateTime: 0n,
startTime: 0n,
isActive: true,
isPaused: false,
pausedAt: null,
Expand All @@ -35,11 +35,11 @@ describe('ClaimableAmountService', () => {

const result = service.getClaimableAmount({
...makeStreamState({
streamId: 1,
streamId: 1n,
ratePerSecond: '5',
depositedAmount: '500',
withdrawnAmount: '100',
lastUpdateTime: 7,
lastUpdateTime: 7n,
}),
});

Expand All @@ -61,7 +61,7 @@ describe('ClaimableAmountService', () => {

const result = service.getClaimableAmount({
...makeStreamState({
streamId: 2,
streamId: 2n,
depositedAmount: '1000',
withdrawnAmount: '900',
}),
Expand All @@ -80,7 +80,7 @@ describe('ClaimableAmountService', () => {

const result = service.getClaimableAmount({
...makeStreamState({
streamId: 3,
streamId: 3n,
withdrawnAmount: '100',
isActive: false,
}),
Expand All @@ -98,7 +98,7 @@ describe('ClaimableAmountService', () => {

const result = service.getClaimableAmount({
...makeStreamState({
streamId: 4,
streamId: 4n,
withdrawnAmount: '150',
}),
});
Expand All @@ -114,7 +114,7 @@ describe('ClaimableAmountService', () => {
});

const input = makeStreamState({
streamId: 5,
streamId: 5n,
ratePerSecond: '7',
depositedAmount: '700',
});
Expand Down Expand Up @@ -146,11 +146,11 @@ describe('ClaimableAmountService', () => {
});

const preWithdrawalState = makeStreamState({
streamId: 7,
streamId: 7n,
ratePerSecond: '10',
depositedAmount: '1000',
withdrawnAmount: '0',
lastUpdateTime: 0,
lastUpdateTime: 0n,
});

// Prime the cache with the pre-withdrawal state.
Expand All @@ -166,11 +166,11 @@ describe('ClaimableAmountService', () => {
// and lastUpdateTime are advanced on the stream row, exactly as
// handleTokensWithdrawn does in soroban-event-worker.ts.
const postWithdrawalState = makeStreamState({
streamId: 7,
streamId: 7n,
ratePerSecond: '10',
depositedAmount: '1000',
withdrawnAmount: '400',
lastUpdateTime: 40,
lastUpdateTime: 40n,
});

// Well within the 60s TTL, so this only passes if the state change (not
Expand All @@ -191,7 +191,7 @@ describe('ClaimableAmountService', () => {

const result = service.getClaimableAmount({
...makeStreamState({
streamId: 6,
streamId: 6n,
ratePerSecond: i128Max,
depositedAmount: i128Max,
withdrawnAmount: '42',
Expand Down Expand Up @@ -228,11 +228,11 @@ describe('ClaimableAmountService', () => {

const result = service.getClaimableAmount({
...makeStreamState({
streamId: 10_000 + iteration,
streamId: BigInt(10_000 + iteration),
ratePerSecond: rate.toString(),
depositedAmount: deposited.toString(),
withdrawnAmount: withdrawn.toString(),
lastUpdateTime: 0,
lastUpdateTime: 0n,
isPaused: paused,
pausedAt: paused ? Number(pauseStart) : null,
totalPausedDuration: paused ? Number(elapsed - pauseStart) : 0,
Expand Down
2 changes: 1 addition & 1 deletion backend/tests/eventRace.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ describe('Action Controller vs Worker Event Write Race Guard (Issue #831)', () =
},
},
create: expect.objectContaining({
streamId: 100,
streamId: BigInt(100),
eventType: 'WITHDRAWN',
transactionHash: 'tx_race_123',
}),
Expand Down
2 changes: 1 addition & 1 deletion backend/tests/events-wire-format.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,7 @@ describe('event wire format', () => {

expect(decodedKeys).toEqual([...allFields].sort());

for (const field of HANDLER_READ_FIELDS[eventName]) {
for (const field of HANDLER_READ_FIELDS[eventName]!) {
expect(decoded).toHaveProperty(field);
}
},
Expand Down
6 changes: 4 additions & 2 deletions backend/tests/indexer-state.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
import { describe, it, expect, vi, beforeEach } from 'vitest';

const mockFindUnique = vi.fn();
const mockCreate = vi.fn();
const { mockFindUnique, mockCreate } = vi.hoisted(() => ({
mockFindUnique: vi.fn(),
mockCreate: vi.fn(),
}));

vi.mock('../src/lib/prisma.js', () => ({
prisma: {
Expand Down
1 change: 1 addition & 0 deletions backend/tests/integration/admin-metrics.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ vi.mock('../../src/lib/redis.js', () => ({
vi.mock('../../src/lib/prisma.js', () => ({
default: mocks.prisma,
prisma: mocks.prisma,
pool: { totalCount: 0, idleCount: 0, waitingCount: 0 },
}));

vi.mock('../../src/middleware/auth.js', async () => {
Expand Down
2 changes: 1 addition & 1 deletion backend/tests/integration/pause-resume.regression.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ describe('Regression #804: Pause/resume controller duplicate StreamEvent', () =>
expect(mockPrisma.stream.update).toHaveBeenCalledTimes(1);
expect(mockPrisma.stream.update).toHaveBeenCalledWith(
expect.objectContaining({
where: { streamId },
where: { streamId: BigInt(77) },
data: expect.objectContaining({ isPaused: true }),
}),
);
Expand Down
5 changes: 3 additions & 2 deletions backend/tests/integration/stream-actions.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ const {
},
streamEvent: {
create: vi.fn(),
upsert: vi.fn(),
findMany: vi.fn().mockResolvedValue([]),
count: vi.fn().mockResolvedValue(0),
},
Expand Down Expand Up @@ -214,9 +215,9 @@ describe('stream action routes', () => {
amount: '100',
});
expect(mockWithdraw).toHaveBeenCalledWith(11n, recipient.publicKey());
expect(mockPrisma.streamEvent.create).toHaveBeenCalledWith(
expect(mockPrisma.streamEvent.upsert).toHaveBeenCalledWith(
expect.objectContaining({
data: expect.objectContaining({
create: expect.objectContaining({
eventType: 'WITHDRAWN',
amount: '100',
transactionHash: 'withdraw-tx-hash',
Expand Down
Loading
Loading