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
50 changes: 50 additions & 0 deletions src/config/configuration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,56 @@ export default () => ({
eventPageLimit: parseInt(process.env.SOROBAN_EVENT_PAGE_LIMIT ?? '100', 10),
},

blockchainIndexer: {
enabled: process.env.BLOCKCHAIN_INDEXER_ENABLED === 'true',
pollIntervalMs: parseInt(
process.env.BLOCKCHAIN_INDEXER_POLL_INTERVAL_MS ?? '2000',
10,
),
maxBackfillLedgers: parseInt(
process.env.BLOCKCHAIN_INDEXER_MAX_BACKFILL_LEDGERS ?? '1000',
10,
),
includeFailed: process.env.BLOCKCHAIN_INDEXER_INCLUDE_FAILED === 'true',
pageLimit: parseInt(process.env.BLOCKCHAIN_INDEXER_PAGE_LIMIT ?? '200', 10),
streamTtlSecs: parseInt(
process.env.BLOCKCHAIN_INDEXER_STREAM_TTL_SECS ?? '300',
10,
),
// Real-time WebSocket streaming settings
wsReconnectBaseDelayMs: parseInt(
process.env.BLOCKCHAIN_INDEXER_WS_RECONNECT_BASE_MS ?? '1000',
10,
),
wsReconnectMaxDelayMs: parseInt(
process.env.BLOCKCHAIN_INDEXER_WS_RECONNECT_MAX_MS ?? '30000',
10,
),
// In-memory event buffer
eventBufferSize: parseInt(
process.env.BLOCKCHAIN_INDEXER_BUFFER_SIZE ?? '10000',
10,
),
eventBufferTtlMs: parseInt(
process.env.BLOCKCHAIN_INDEXER_BUFFER_TTL_MS ?? '60000',
10,
),
// Batched persistence
batchFlushIntervalMs: parseInt(
process.env.BLOCKCHAIN_INDEXER_BATCH_FLUSH_MS ?? '1000',
10,
),
batchMaxSize: parseInt(
process.env.BLOCKCHAIN_INDEXER_BATCH_MAX_SIZE ?? '1000',
10,
),
// Retention policy
retentionDays: parseInt(
process.env.BLOCKCHAIN_INDEXER_RETENTION_DAYS ?? '90',
10,
),
},

rateLimit: {
strategy: process.env.RATE_LIMIT_STRATEGY ?? 'sliding_window',
defaultWindowSize: parseInt(process.env.RATE_LIMIT_DEFAULT_WINDOW ?? '60', 10),
Expand Down
8 changes: 8 additions & 0 deletions src/config/env.validation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,14 @@ export const envValidationSchema = Joi.object({
BLOCKCHAIN_INDEXER_MAX_BACKFILL_LEDGERS: Joi.number().default(1000),
BLOCKCHAIN_INDEXER_STREAM_TTL_SECS: Joi.number().default(300),
BLOCKCHAIN_INDEXER_INCLUDE_FAILED: Joi.boolean().default(true),
BLOCKCHAIN_INDEXER_PAGE_LIMIT: Joi.number().default(200),
BLOCKCHAIN_INDEXER_WS_RECONNECT_BASE_MS: Joi.number().default(1000),
BLOCKCHAIN_INDEXER_WS_RECONNECT_MAX_MS: Joi.number().default(30000),
BLOCKCHAIN_INDEXER_BUFFER_SIZE: Joi.number().default(10000),
BLOCKCHAIN_INDEXER_BUFFER_TTL_MS: Joi.number().default(60000),
BLOCKCHAIN_INDEXER_BATCH_FLUSH_MS: Joi.number().default(1000),
BLOCKCHAIN_INDEXER_BATCH_MAX_SIZE: Joi.number().default(1000),
BLOCKCHAIN_INDEXER_RETENTION_DAYS: Joi.number().default(90),

// Rate limiting
RATE_LIMIT_STRATEGY: Joi.string()
Expand Down
79 changes: 79 additions & 0 deletions src/modules/blockchain-indexer/blockchain-indexer.controller.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,13 @@ import {
} from '@nestjs/swagger';
import { JwtAuthGuard } from '../auth/guards/jwt-auth.guard';
import { BlockchainIndexerService } from './blockchain-indexer.service';
import { LedgerIndexerService } from './services/ledger-indexer.service';
import { EventQueryService } from './services/event-query.service';
import { QueryEventsDto } from './dto/query-events.dto';
import { TemporalQueryDto } from './dto/subscribe-events.dto';
import { PaginatedResultDto } from '@app/common';
import { BlockchainEvent } from './entities/blockchain-event.entity';
import { IndexedEvent } from './entities/indexed-event.entity';
import { Response } from 'express';

@ApiTags('blockchain-indexer')
Expand All @@ -30,9 +33,12 @@ import { Response } from 'express';
export class BlockchainIndexerController {
constructor(
private readonly blockchainIndexerService: BlockchainIndexerService,
private readonly ledgerIndexerService: LedgerIndexerService,
private readonly queryService: EventQueryService,
) {}

// ─── Legacy event queries ───────────────────────────────────────

@Get('events')
@ApiOperation({ summary: 'Query indexed blockchain events with filtering' })
@ApiResponse({
Expand Down Expand Up @@ -67,6 +73,77 @@ export class BlockchainIndexerController {
return this.queryService.findByTransactionHash(transactionHash);
}

// ─── Real-time indexer endpoints ────────────────────────────────

@Get('realtime/status')
@ApiOperation({
summary: 'Get real-time ledger indexer status, metrics, and health',
})
@ApiResponse({
status: HttpStatus.OK,
description: 'Real-time indexer status',
})
async getRealtimeStatus() {
return this.ledgerIndexerService.getStatus();
}

@Get('realtime/events')
@ApiOperation({
summary:
'Query indexed events from the real-time pipeline with filtering',
})
@ApiResponse({
status: HttpStatus.OK,
description: 'Events retrieved successfully',
})
async getRealtimeEvents(@Query() queryDto: QueryEventsDto) {
const result = await this.queryService.findEvents({
...queryDto,
skip: queryDto.skip,
startTime: queryDto.startTime ? new Date(queryDto.startTime) : undefined,
endTime: queryDto.endTime ? new Date(queryDto.endTime) : undefined,
excludeInvalidated: true,
});
return new PaginatedResultDto<BlockchainEvent>(
result.events,
result.total,
queryDto.page,
queryDto.limit,
);
}

@Get('realtime/temporal')
@ApiOperation({
summary: 'Query state at a specific block (temporal query)',
description:
'Returns events as they were at a given ledger sequence, supporting ' +
'historical state reconstruction.',
})
@ApiResponse({
status: HttpStatus.OK,
description: 'Temporal query results',
})
async getTemporalState(@Query() queryDto: TemporalQueryDto) {
const result = await this.queryService.findEvents({
ledgerFrom: 0,
ledgerTo: queryDto.atLedger,
eventType: queryDto.eventType,
sourceAccount: queryDto.account,
skip: queryDto.skip,
limit: queryDto.limit,
excludeInvalidated: true,
});
return {
asOfLedger: queryDto.atLedger,
...new PaginatedResultDto<IndexedEvent>(
result.events as any,
result.total,
queryDto.page,
queryDto.limit,
),
};
}

@Get('events/stream')
@ApiOperation({ summary: 'Stream recent blockchain events via SSE' })
@ApiResponse({
Expand Down Expand Up @@ -96,6 +173,8 @@ export class BlockchainIndexerController {
});
}

// ─── Health and operational endpoints ───────────────────────────

@Get('status')
@ApiOperation({ summary: 'Get blockchain indexer status and metrics' })
@ApiResponse({
Expand Down
26 changes: 24 additions & 2 deletions src/modules/blockchain-indexer/blockchain-indexer.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,30 +7,52 @@ import { BlockchainIndexerController } from './blockchain-indexer.controller';
import { BlockchainIndexerService } from './blockchain-indexer.service';
import { BlockchainEvent } from './entities/blockchain-event.entity';
import { IndexerState } from './entities/indexer-state.entity';
import { IndexedEvent } from './entities/indexed-event.entity';
import { StellarEventSourceService } from './services/stellar-event-source.service';
import { IndexingStateService } from './services/indexing-state.service';
import { ReorgHandlerService } from './services/reorg-handler.service';
import { EventIndexerService } from './services/event-indexer.service';
import { EventQueryService } from './services/event-query.service';
import { EventStreamService } from './services/event-stream.service';
// New real-time indexing services
import { HorizonStreamService } from './services/horizon-stream.service';
import { EventNormalizer } from './services/event-normalizer.service';
import { EventBufferService } from './services/event-buffer.service';
import { BatchedPersistenceService } from './services/batched-persistence.service';
import { SubscriptionManager } from './services/subscription-manager.service';
import { LedgerIndexerService } from './services/ledger-indexer.service';
import { EventWebSocketGateway } from './event-websocket.gateway';

@Module({
imports: [
ConfigModule,
TypeOrmModule.forFeature([BlockchainEvent, IndexerState]),
TypeOrmModule.forFeature([BlockchainEvent, IndexerState, IndexedEvent]),
StellarModule,
RedisModule,
],
controllers: [BlockchainIndexerController],
providers: [
// Legacy polling-based indexer (kept for backward compatibility)
BlockchainIndexerService,
StellarEventSourceService,
IndexingStateService,
ReorgHandlerService,
EventIndexerService,
EventQueryService,
EventStreamService,
// Real-time streaming indexer
HorizonStreamService,
EventNormalizer,
EventBufferService,
BatchedPersistenceService,
SubscriptionManager,
LedgerIndexerService,
EventWebSocketGateway,
],
exports: [
BlockchainIndexerService,
LedgerIndexerService,
SubscriptionManager,
],
exports: [BlockchainIndexerService],
})
export class BlockchainIndexerModule {}
100 changes: 100 additions & 0 deletions src/modules/blockchain-indexer/dto/subscribe-events.dto.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
import { ApiPropertyOptional } from '@nestjs/swagger';
import { Type } from 'class-transformer';
import {
IsArray,
IsEnum,
IsInt,
IsOptional,
IsString,
Max,
Min,
} from 'class-validator';
import { BlockchainEventType } from '../enums/blockchain-event-type.enum';

/**
* DTO for subscribing to real-time events via REST.
* The WebSocket gateway accepts a similar shape directly.
*/
export class SubscribeEventsDto {
@ApiPropertyOptional({
description: 'Filter by event types. Empty = all types.',
enum: BlockchainEventType,
isArray: true,
})
@IsOptional()
@IsArray()
@IsEnum(BlockchainEventType, { each: true })
eventTypes?: BlockchainEventType[];

@ApiPropertyOptional({
description: 'Filter by Soroban contract IDs.',
type: [String],
})
@IsOptional()
@IsArray()
@IsString({ each: true })
contractIds?: string[];

@ApiPropertyOptional({
description: 'Filter by source or destination accounts.',
type: [String],
})
@IsOptional()
@IsArray()
@IsString({ each: true })
accounts?: string[];

@ApiPropertyOptional({ description: 'Minimum ledger sequence to include.' })
@IsOptional()
@Type(() => Number)
@IsInt()
@Min(0)
fromLedger?: number;
}

/**
* DTO for querying state at a specific block (temporal query).
*/
export class TemporalQueryDto {
@ApiPropertyOptional({
description: 'Query events as they were at this ledger sequence.',
})
@Type(() => Number)
@IsInt()
@Min(0)
atLedger: number;

@ApiPropertyOptional({ description: 'Filter by event types.' })
@IsOptional()
@IsEnum(BlockchainEventType)
eventType?: BlockchainEventType;

@ApiPropertyOptional({ description: 'Filter by contract ID.' })
@IsOptional()
@IsString()
contractId?: string;

@ApiPropertyOptional({ description: 'Filter by account.' })
@IsOptional()
@IsString()
account?: string;

@ApiPropertyOptional({ default: 1, minimum: 1 })
@IsOptional()
@Type(() => Number)
@IsInt()
@Min(1)
page = 1;

@ApiPropertyOptional({ default: 20, minimum: 1, maximum: 100 })
@IsOptional()
@Type(() => Number)
@IsInt()
@Min(1)
@Max(100)
limit = 20;

get skip(): number {
return (this.page - 1) * this.limit;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,11 @@ export enum BlockchainEventType {
CREATE_ACCOUNT = 'create_account',
ACCOUNT_MERGE = 'account_merge',
TRANSACTION = 'transaction',
// Soroban contract event types
SOROBAN_CONTRACT_INVOCATION = 'soroban_contract_invocation',
SOROBAN_CONTRACT_EVENT = 'soroban_contract_event',
SOROBAN_SYSTEM_EVENT = 'soroban_system_event',
SOROBAN_DIAGNOSTIC_EVENT = 'soroban_diagnostic_event',
}

@Entity('blockchain_events')
Expand Down
Loading
Loading