Skip to content

Commit b1e15d3

Browse files
committed
feat(backend): per-network ledger cursor, gap alerts, reindex reset
- Add ledger_cursors (network PK, last_processed_ledger, updated_at) and ledger_gap_alert_dedup for cooldown-based gap notifications - Indexer advances cursor inside same Prisma tx as raw_events + projections - Gap alert when latestLedger - lastProcessed > threshold; dedupe via cooldown - Admin POST /admin/reindex resets cursor (fromLedger-1), enqueues Bull job; ReindexWorkerService runs processUntilCaughtUp (disabled in test / DISABLE_REINDEX_WORKER) - STELLAR_NETWORK + INDEXER_GAP_ALERT_* env; ClaimsService reads cursor by network - Migration seeds testnet row from legacy indexer_state - Tests: cursor resume, tx-bound advance, gap dedupe, admin reindex - Fix idempotency middleware Redis import path + fail-open on lookup errors - Jest: include *.spec.ts in testMatch Made-with: Cursor
1 parent 3c90de3 commit b1e15d3

17 files changed

Lines changed: 567 additions & 76 deletions

backend/jest.config.js

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
module.exports = {
33
testEnvironment: "node",
44
roots: ["<rootDir>/src", "<rootDir>/tests"],
5-
testMatch: ["**/__tests__/**/*.test.ts", "**/*.test.ts"],
5+
testMatch: ["**/__tests__/**/*.test.ts", "**/*.test.ts", "**/*.spec.ts"],
66
transform: {
77
"^.+\\.tsx?$": ["ts-jest", { tsconfig: "tsconfig.test.json" }],
88
},
Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
-- CreateTable
2+
CREATE TABLE "ledger_cursors" (
3+
"network" TEXT NOT NULL,
4+
"last_processed_ledger" INTEGER NOT NULL,
5+
"updated_at" TIMESTAMP(3) NOT NULL,
6+
7+
CONSTRAINT "ledger_cursors_pkey" PRIMARY KEY ("network")
8+
);
9+
10+
-- CreateTable
11+
CREATE TABLE "ledger_gap_alert_dedup" (
12+
"network" TEXT NOT NULL,
13+
"last_fired_at" TIMESTAMP(3) NOT NULL,
14+
"last_gap_size" INTEGER,
15+
"last_processed_ledger" INTEGER,
16+
"latest_ledger" INTEGER,
17+
18+
CONSTRAINT "ledger_gap_alert_dedup_pkey" PRIMARY KEY ("network")
19+
);
20+
21+
-- Seed default network from legacy indexer_state (if present)
22+
INSERT INTO "ledger_cursors" ("network", "last_processed_ledger", "updated_at")
23+
SELECT 'testnet', COALESCE((SELECT MAX("last_ledger") FROM "indexer_state"), 0), NOW()
24+
WHERE NOT EXISTS (SELECT 1 FROM "ledger_cursors" WHERE "network" = 'testnet');
Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,3 @@
1+
# Please do not edit this file manually
2+
# It should be added in your version-control system (i.e. Git)
3+
provider = "postgresql"

backend/prisma/schema.prisma

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -99,6 +99,28 @@ model IndexerState {
9999
@@map("indexer_state")
100100
}
101101

102+
/// Per-Stellar-network indexer progress. `last_processed_ledger` is the last ledger
103+
/// sequence whose relevant contract events have been applied (cursor advances in the
104+
/// same DB transaction as raw event / projection upserts).
105+
model LedgerCursor {
106+
network String @id
107+
lastProcessedLedger Int @map("last_processed_ledger")
108+
updatedAt DateTime @updatedAt @map("updated_at")
109+
110+
@@map("ledger_cursors")
111+
}
112+
113+
/// One row per network: last time a ledger gap alert was emitted (dedup / cooldown).
114+
model LedgerGapAlertDedup {
115+
network String @id
116+
lastFiredAt DateTime @map("last_fired_at")
117+
lastGapSize Int? @map("last_gap_size")
118+
lastProcessedLedger Int? @map("last_processed_ledger")
119+
latestLedger Int? @map("latest_ledger")
120+
121+
@@map("ledger_gap_alert_dedup")
122+
}
123+
102124
model RawEvent {
103125
id Int @id @default(autoincrement())
104126
txHash String

backend/src/admin/admin.controller.spec.ts

Lines changed: 30 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,20 @@
11
import { Test, TestingModule } from '@nestjs/testing';
22
import { ExecutionContext, ForbiddenException } from '@nestjs/common';
3+
import { ConfigService } from '@nestjs/config';
34
import { Request } from 'express';
45
import { AdminController } from './admin.controller';
56
import { AdminService } from './admin.service';
67
import { AuditService } from './audit.service';
78
import { AdminRoleGuard } from './guards/admin-role.guard';
89
import { JwtAuthGuard } from '../auth/guards/jwt-auth.guard';
10+
import { PrivacyService } from '../maintenance/privacy.service';
11+
import { RateLimitService } from '../rate-limit/rate-limit.service';
912

1013
const mockAdminService = { enqueueReindex: jest.fn(), setFeatureFlag: jest.fn(), getFeatureFlags: jest.fn() };
1114
const mockAuditService = { write: jest.fn(), findAll: jest.fn() };
15+
const mockConfigService = {
16+
get: jest.fn((key: string, def?: string) => (key === 'STELLAR_NETWORK' ? 'testnet' : def)),
17+
};
1218

1319
const adminReq = (role = 'admin') => ({ user: { walletAddress: 'GADMIN', role }, ip: '127.0.0.1' });
1420
const toExecutionContext = (role?: string): ExecutionContext =>
@@ -26,6 +32,9 @@ describe('AdminController', () => {
2632
providers: [
2733
{ provide: AdminService, useValue: mockAdminService },
2834
{ provide: AuditService, useValue: mockAuditService },
35+
{ provide: ConfigService, useValue: mockConfigService },
36+
{ provide: PrivacyService, useValue: {} },
37+
{ provide: RateLimitService, useValue: {} },
2938
],
3039
})
3140
.overrideGuard(JwtAuthGuard).useValue({ canActivate: () => true })
@@ -43,12 +52,30 @@ describe('AdminController', () => {
4352
it('enqueues job and writes audit row', async () => {
4453
mockAdminService.enqueueReindex.mockResolvedValue('job-123');
4554
const result = await controller.reindex({ fromLedger: 500 }, adminReq() as unknown as Request);
46-
expect(result).toEqual({ jobId: 'job-123', fromLedger: 500, status: 'queued' });
47-
expect(mockAdminService.enqueueReindex).toHaveBeenCalledWith(500);
55+
expect(result).toEqual({
56+
jobId: 'job-123',
57+
fromLedger: 500,
58+
network: 'testnet',
59+
status: 'queued',
60+
});
61+
expect(mockAdminService.enqueueReindex).toHaveBeenCalledWith(500, 'testnet');
4862
expect(mockAuditService.write).toHaveBeenCalledWith(
49-
expect.objectContaining({ actor: 'GADMIN', action: 'reindex', payload: expect.objectContaining({ fromLedger: 500 }) }),
63+
expect.objectContaining({
64+
actor: 'GADMIN',
65+
action: 'reindex',
66+
payload: expect.objectContaining({ fromLedger: 500, network: 'testnet' }),
67+
}),
5068
);
5169
});
70+
71+
it('passes explicit network to enqueue', async () => {
72+
mockAdminService.enqueueReindex.mockResolvedValue('job-456');
73+
await controller.reindex(
74+
{ fromLedger: 100, network: 'public' },
75+
adminReq() as unknown as Request,
76+
);
77+
expect(mockAdminService.enqueueReindex).toHaveBeenCalledWith(100, 'public');
78+
});
5279
});
5380

5481
describe('GET /admin/audits', () => {

backend/src/admin/admin.controller.ts

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import {
1212
HttpCode,
1313
HttpStatus,
1414
} from '@nestjs/common';
15+
import { ConfigService } from '@nestjs/config';
1516
import { ApiBearerAuth, ApiOperation, ApiTags } from '@nestjs/swagger';
1617
import { IsEnum, IsOptional, IsString } from 'class-validator';
1718
import { Request } from 'express';
@@ -48,6 +49,7 @@ export class AdminController {
4849
private readonly auditService: AuditService,
4950
private readonly privacyService: PrivacyService,
5051
private readonly rateLimitService: RateLimitService,
52+
private readonly configService: ConfigService,
5153
) {}
5254

5355
/**
@@ -64,14 +66,16 @@ export class AdminController {
6466
@ApiOperation({ summary: 'Enqueue a ledger reindex job from a given ledger' })
6567
async reindex(@Body() dto: ReindexDto, @Req() req: AdminRequest) {
6668
const actor = req.user?.walletAddress ?? 'unknown';
67-
const jobId = await this.adminService.enqueueReindex(dto.fromLedger);
69+
const network =
70+
dto.network ?? this.configService.get<string>('STELLAR_NETWORK', 'testnet');
71+
const jobId = await this.adminService.enqueueReindex(dto.fromLedger, network);
6872
await this.auditService.write({
6973
actor,
7074
action: 'reindex',
71-
payload: { fromLedger: dto.fromLedger, jobId },
75+
payload: { fromLedger: dto.fromLedger, network, jobId },
7276
ipAddress: req.ip,
7377
});
74-
return { jobId, fromLedger: dto.fromLedger, status: 'queued' };
78+
return { jobId, fromLedger: dto.fromLedger, network, status: 'queued' };
7579
}
7680

7781
/**

backend/src/admin/admin.module.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { Module } from '@nestjs/common';
2+
import { ConfigModule } from '@nestjs/config';
23
import { AdminController } from './admin.controller';
34
import { AdminService } from './admin.service';
45
import { AuditService } from './audit.service';
@@ -8,7 +9,7 @@ import { MaintenanceModule } from '../maintenance/maintenance.module';
89
import { RateLimitModule } from '../rate-limit/rate-limit.module';
910

1011
@Module({
11-
imports: [PrismaModule, AuthModule, MaintenanceModule, RateLimitModule],
12+
imports: [ConfigModule, PrismaModule, AuthModule, MaintenanceModule, RateLimitModule],
1213
controllers: [AdminController],
1314
providers: [AdminService, AuditService],
1415
exports: [AuditService],
Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,62 @@
1+
const mockQueueAdd = jest.fn().mockResolvedValue({ id: 'queued-job-id' });
2+
3+
jest.mock('bullmq', () => ({
4+
Queue: jest.fn().mockImplementation(() => ({
5+
add: (...args: unknown[]) => mockQueueAdd(...args),
6+
})),
7+
}));
8+
9+
jest.mock('../redis/client', () => ({
10+
getBullMQConnection: () => ({}),
11+
}));
12+
13+
import { AdminService } from './admin.service';
14+
15+
describe('AdminService', () => {
16+
beforeEach(() => {
17+
jest.clearAllMocks();
18+
mockQueueAdd.mockResolvedValue({ id: 'queued-job-id' });
19+
});
20+
21+
describe('enqueueReindex', () => {
22+
it('sets last_processed_ledger to fromLedger-1 and enqueues with network', async () => {
23+
const upsert = jest.fn();
24+
const prisma = {
25+
$transaction: jest.fn(async (fn: (t: { ledgerCursor: { upsert: jest.Mock } }) => Promise<void>) =>
26+
fn({ ledgerCursor: { upsert } })),
27+
};
28+
29+
const svc = new AdminService(prisma as never);
30+
const jobId = await svc.enqueueReindex(500, 'testnet');
31+
32+
expect(jobId).toBe('queued-job-id');
33+
expect(upsert).toHaveBeenCalledWith({
34+
where: { network: 'testnet' },
35+
create: { network: 'testnet', lastProcessedLedger: 499 },
36+
update: { lastProcessedLedger: 499 },
37+
});
38+
expect(mockQueueAdd).toHaveBeenCalledWith(
39+
'reindex',
40+
{ fromLedger: 500, network: 'testnet' },
41+
expect.objectContaining({
42+
jobId: expect.stringMatching(/^reindex-testnet-500-/),
43+
}),
44+
);
45+
});
46+
47+
it('clamps at 0 when fromLedger is 0', async () => {
48+
const upsert = jest.fn();
49+
const prisma = {
50+
$transaction: jest.fn(async (fn: (t: { ledgerCursor: { upsert: jest.Mock } }) => Promise<void>) =>
51+
fn({ ledgerCursor: { upsert } })),
52+
};
53+
const svc = new AdminService(prisma as never);
54+
await svc.enqueueReindex(0, 'public');
55+
expect(upsert).toHaveBeenCalledWith(
56+
expect.objectContaining({
57+
create: expect.objectContaining({ lastProcessedLedger: 0 }),
58+
}),
59+
);
60+
});
61+
});
62+
});

backend/src/admin/admin.service.ts

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,9 +20,25 @@ export class AdminService {
2020
});
2121
}
2222

23-
async enqueueReindex(fromLedger: number): Promise<string> {
24-
const job = await this.reindexQueue.add('reindex', { fromLedger }, { jobId: `reindex-${fromLedger}-${Date.now()}` });
25-
this.logger.log(`Reindex job enqueued: ${job.id} from ledger ${fromLedger}`);
23+
/**
24+
* Reset per-network cursor so the next indexer pass starts at `fromLedger`,
25+
* then enqueue a BullMQ job to drive catch-up (see ReindexWorkerService).
26+
*/
27+
async enqueueReindex(fromLedger: number, network: string): Promise<string> {
28+
const lastProcessed = Math.max(0, fromLedger - 1);
29+
await this.prisma.$transaction(async (tx) => {
30+
await tx.ledgerCursor.upsert({
31+
where: { network },
32+
create: { network, lastProcessedLedger: lastProcessed },
33+
update: { lastProcessedLedger: lastProcessed },
34+
});
35+
});
36+
const job = await this.reindexQueue.add(
37+
'reindex',
38+
{ fromLedger, network },
39+
{ jobId: `reindex-${network}-${fromLedger}-${Date.now()}` },
40+
);
41+
this.logger.log(`Reindex job enqueued: ${job.id} network=${network} fromLedger=${fromLedger}`);
2642
return job.id!;
2743
}
2844

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,19 @@
1-
import { IsInt, Min } from 'class-validator';
2-
import { ApiProperty } from '@nestjs/swagger';
1+
import { IsInt, IsOptional, IsString, Matches, Min } from 'class-validator';
2+
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
33

44
export class ReindexDto {
55
@ApiProperty({ description: 'Ledger sequence to reindex from', minimum: 0 })
66
@IsInt()
77
@Min(0)
88
fromLedger!: number;
9+
10+
@ApiPropertyOptional({
11+
description:
12+
'Stellar logical network id (must match STELLAR_NETWORK / indexer cursor row). Defaults to server config.',
13+
example: 'testnet',
14+
})
15+
@IsOptional()
16+
@IsString()
17+
@Matches(/^[a-z0-9][a-z0-9_-]{0,62}$/i)
18+
network?: string;
919
}

0 commit comments

Comments
 (0)