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
10 changes: 10 additions & 0 deletions backend/src/api/rest/monitor/monitor.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,17 @@ import { ConfigModule } from '@nestjs/config';
import { SupabaseModule } from '../../../services/supabase.module';
import { MonitorController } from './monitor.controller';
import { MonitorService } from './monitor.service';
import { ReplayController } from './replay.controller';
import { ReplayService } from '../../../services/replay.service';
import { AdminGuard } from './admin.guard';
import { BackfillLock } from '../../../services/backfill-lock';
import { HorizonClientService } from '../../../services/horizon-client.service';
import { IndexerService } from '../../../services/indexer.service';

@Module({
imports: [SupabaseModule, ConfigModule],
controllers: [MonitorController, ReplayController],
providers: [MonitorService, ReplayService, AdminGuard, BackfillLock, HorizonClientService, IndexerService],
import { AuditLogInterceptor } from './audit-log.interceptor';

@Module({
Expand Down
109 changes: 109 additions & 0 deletions backend/src/api/rest/monitor/replay.controller.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
import {
Controller,
Post,
Get,
Body,
Param,
UseGuards,
HttpCode,
HttpStatus,
BadRequestException,
NotFoundException,
} from '@nestjs/common';
import { ApiTags, ApiOperation, ApiResponse, ApiSecurity } from '@nestjs/swagger';
import { ReplayService, type ReplayJobConfig, type ReplayJobStatus } from '../../../services/replay.service';
import { AdminGuard } from './admin.guard';

@ApiTags('Admin - Replay')
@ApiSecurity('admin-token')
@UseGuards(AdminGuard)
@Controller('admin/replay')
export class ReplayController {
constructor(private readonly replayService: ReplayService) {}

/**
* POST /admin/replay — Start a new replay job
* Requires admin token in X-Admin-Token header
*/
@Post()
@HttpCode(HttpStatus.ACCEPTED)
@ApiOperation({
summary: 'Start a ledger replay job',
description: 'Triggers an async replay of ledgers in the specified range. Returns a job ID for polling progress.',
})
@ApiResponse({
status: 202,
description: 'Replay job started',
schema: {
type: 'object',
properties: {
jobId: { type: 'string', format: 'uuid' },
message: { type: 'string' },
},
},
})
@ApiResponse({
status: 400,
description: 'Invalid replay configuration',
})
@ApiResponse({
status: 401,
description: 'Invalid or missing admin token',
})
async startReplay(@Body() config: ReplayJobConfig) {
try {
const jobId = this.replayService.startReplay(config);
return {
jobId,
message: `Replay job started. Poll /admin/replay/${jobId} for progress.`,
};
} catch (err) {
throw new BadRequestException(
err instanceof Error ? err.message : 'Invalid replay configuration',
);
}
}

/**
* GET /admin/replay/:jobId — Get replay job status
* Requires admin token in X-Admin-Token header
*/
@Get(':jobId')
@ApiOperation({
summary: 'Get replay job status',
description: 'Poll the status and progress of a replay job.',
})
@ApiResponse({
status: 200,
description: 'Job status',
schema: {
type: 'object',
properties: {
jobId: { type: 'string', format: 'uuid' },
status: { type: 'string', enum: ['pending', 'running', 'completed', 'failed'] },
config: { type: 'object' },
progress: { type: 'object' },
result: { type: 'object' },
error: { type: 'string' },
createdAt: { type: 'string', format: 'date-time' },
startedAt: { type: 'string', format: 'date-time' },
completedAt: { type: 'string', format: 'date-time' },
},
},
})
@ApiResponse({
status: 404,
description: 'Job not found',
})
@ApiResponse({
status: 401,
description: 'Invalid or missing admin token',
})
getJobStatus(@Param('jobId') jobId: string): ReplayJobStatus {
const job = this.replayService.getJobStatus(jobId);
if (!job) {
throw new NotFoundException(`Replay job ${jobId} not found`);
}
return job;
}
}
228 changes: 228 additions & 0 deletions backend/src/services/replay.service.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,228 @@
import { Injectable, Logger } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { randomUUID } from 'crypto';
import { BackfillLock, BackfillLockError } from './backfill-lock';
import { HorizonClientService } from './horizon-client.service';
import { IndexerService } from './indexer.service';

export interface ReplayJobConfig {
fromLedger: number;
toLedger: number;
contractId?: string;
dryRun?: boolean;
}

export interface ReplayJobStatus {
jobId: string;
status: 'pending' | 'running' | 'completed' | 'failed';
config: ReplayJobConfig;
progress: {
processedCount: number;
skippedCount: number;
totalLedgers: number;
currentLedger?: number;
};
result?: {
elapsedMs: number;
missingLedgers: number[];
};
error?: string;
createdAt: string;
startedAt?: string;
completedAt?: string;
}

@Injectable()
export class ReplayService {
private readonly logger = new Logger(ReplayService.name);
private readonly jobs = new Map<string, ReplayJobStatus>();
private readonly maxRange: number;
private readonly retryCount: number;
private readonly retryDelayMs: number;

constructor(
private readonly config: ConfigService,
private readonly lock: BackfillLock,
private readonly horizonClient: HorizonClientService,
private readonly indexer: IndexerService,
) {
this.maxRange = this.config.get<number>('BACKFILL_MAX_RANGE', 10000);
this.retryCount = this.config.get<number>('BACKFILL_RETRY_COUNT', 3);
this.retryDelayMs = this.config.get<number>('BACKFILL_RETRY_DELAY_MS', 1000);
}

/**
* Start a new replay job
*/
startReplay(config: ReplayJobConfig): string {
this.validateConfig(config);

const jobId = randomUUID();
const job: ReplayJobStatus = {
jobId,
status: 'pending',
config,
progress: {
processedCount: 0,
skippedCount: 0,
totalLedgers: config.toLedger - config.fromLedger + 1,
},
createdAt: new Date().toISOString(),
};

this.jobs.set(jobId, job);

// Start async replay in background
this.runReplay(jobId).catch((err) => {
this.logger.error(`Replay job ${jobId} failed:`, err);
const job = this.jobs.get(jobId);
if (job) {
job.status = 'failed';
job.error = err instanceof Error ? err.message : String(err);
job.completedAt = new Date().toISOString();
}
});

return jobId;
}

/**
* Get replay job status
*/
getJobStatus(jobId: string): ReplayJobStatus | null {
return this.jobs.get(jobId) || null;
}

/**
* Validate replay configuration
*/
private validateConfig(config: ReplayJobConfig): void {
if (!Number.isInteger(config.fromLedger) || config.fromLedger <= 0) {
throw new Error(
`fromLedger must be a positive integer, got: ${config.fromLedger}`,
);
}

if (!Number.isInteger(config.toLedger) || config.toLedger <= 0) {
throw new Error(
`toLedger must be a positive integer, got: ${config.toLedger}`,
);
}

if (config.fromLedger > config.toLedger) {
throw new Error(
`fromLedger (${config.fromLedger}) must be <= toLedger (${config.toLedger})`,
);
}

const range = config.toLedger - config.fromLedger + 1;
if (range > this.maxRange) {
throw new Error(
`Range of ${range} ledgers exceeds BACKFILL_MAX_RANGE (${this.maxRange})`,
);
}
}

/**
* Run the replay job
*/
private async runReplay(jobId: string): Promise<void> {
const job = this.jobs.get(jobId);
if (!job) throw new Error(`Job ${jobId} not found`);

job.status = 'running';
job.startedAt = new Date().toISOString();

// Acquire lock to prevent concurrent indexing
if (!this.lock.tryAcquire()) {
throw new BackfillLockError('Could not acquire backfill lock');
}

try {
const { fromLedger, toLedger, dryRun } = job.config;
const missingLedgers: number[] = [];
const runStart = Date.now();

this.logger.log(
`Replay job ${jobId} started: fromLedger=${fromLedger} toLedger=${toLedger} dryRun=${dryRun}`,
);

for (let seq = fromLedger; seq <= toLedger; seq++) {
job.progress.currentLedger = seq;

let ledgerData: Awaited<ReturnType<HorizonClientService['fetchLedger']>> = null;
let lastError: Error | null = null;
let attempts = 0;

// Fetch with retry
while (attempts <= this.retryCount) {
try {
ledgerData = await this.horizonClient.fetchLedger(seq);
lastError = null;
break;
} catch (err) {
lastError = err instanceof Error ? err : new Error(String(err));
attempts++;
if (attempts <= this.retryCount) {
await this.delay(this.retryDelayMs);
}
}
}

if (lastError !== null) {
this.logger.warn(
`Replay job ${jobId}: Missing ledger seq=${seq}: transient error after ${attempts} attempt(s)`,
);
missingLedgers.push(seq);
job.progress.skippedCount++;
} else if (ledgerData === null) {
this.logger.warn(
`Replay job ${jobId}: Missing ledger seq=${seq}: not found in Horizon archive`,
);
missingLedgers.push(seq);
job.progress.skippedCount++;
} else {
// Submit to indexer (unless dry-run)
if (!dryRun) {
await this.indexer.submitLedger(ledgerData, seq);
} else {
this.logger.debug(
`Replay job ${jobId} (dry-run): Would submit ledger seq=${seq}`,
);
}
job.progress.processedCount++;
}

// Progress log every 100 ledgers
const processed = seq - fromLedger + 1;
const total = job.progress.totalLedgers;
if (processed % 100 === 0) {
const pct = Math.round((processed / total) * 100);
this.logger.log(
`Replay job ${jobId} progress: seq=${seq} processed=${processed}/${total} (${pct}%)`,
);
}
}

const elapsedMs = Date.now() - runStart;

job.status = 'completed';
job.completedAt = new Date().toISOString();
job.result = {
elapsedMs,
missingLedgers,
};

this.logger.log(
`Replay job ${jobId} completed: processedCount=${job.progress.processedCount} skippedCount=${job.progress.skippedCount} elapsedMs=${elapsedMs}`,
);
} finally {
this.lock.release();
delete job.progress.currentLedger;
}
}

private delay(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
}
Loading