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
6 changes: 6 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,12 @@ ADMIN_SECRET=your_admin_secret_here
# Format: contractId=name separated by commas.
# CONTRACT_NAMES=your_settlement_contract_id=settlement,your_governance_contract_id=governance

# Indexer — smart startup ledger discovery (#352)
# Number of ledgers to backfill from the network tip on fresh deployments (default: 1000).
# INITIAL_BACKFILL_LEDGERS=1000
# Manual override: skip auto-discovery and start from this specific ledger.
# INDEX_FROM_LEDGER=5000

# Redis (required by settlement engine BullMQ worker)
REDIS_URL=redis://localhost:6379
# BullMQ Redis connection tuning. BullMQ expects enableReadyCheck=false.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
-- DropIndex
DROP INDEX "IndexedEvent_stellarId_key";

-- CreateIndex
CREATE UNIQUE INDEX "IndexedEvent_stellarId_contractId_ledger_key" ON "IndexedEvent"("stellarId", "contractId", "ledger");
3 changes: 2 additions & 1 deletion prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ model Settlement {

model IndexedEvent {
id String @id
stellarId String? @unique
stellarId String?
contractId String
contractName String?
topics String[]
Expand All @@ -96,6 +96,7 @@ model IndexedEvent {
ledger Int
indexedAt DateTime @default(now())

@@unique([stellarId, contractId, ledger])
@@index([ledger])
@@index([type])
@@index([contractId])
Expand Down
2 changes: 1 addition & 1 deletion services/indexer/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
"start": "node dist/index.js",
"type-check": "tsc --noEmit",
"clean": "rm -rf dist",
"test": "TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/worker-shutdown.test.ts && TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/cleanup.test.ts && TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/webhook-cache.test.ts && cross-env DOTENV_CONFIG_PATH=../../.env TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/index.test.ts"
"test": "TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/worker-shutdown.test.ts && TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/cleanup.test.ts && TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/webhook-cache.test.ts && TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/startup.test.ts && cross-env DOTENV_CONFIG_PATH=../../.env TS_NODE_TRANSPILE_ONLY=true node --loader ts-node/esm src/index.test.ts"
},
"dependencies": {
"@bettapay/webhook-delivery": "workspace:*",
Expand Down
235 changes: 201 additions & 34 deletions services/indexer/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ import Fastify from 'fastify';
import cors from '@fastify/cors';
import crypto from 'crypto';
import { Queue, Worker } from 'bullmq';
import { createWebhookQueue, createWebhookWorker } from '@bettapay/webhook-delivery';
import { createWebhookQueue, createWebhookWorker, WEBHOOK_DEFAULTS } from '@bettapay/webhook-delivery';
import { closeWorkerWithTimeout, trackActiveJob } from './worker-shutdown.js';
import { PrismaClient, WebhookSubscription } from '@prisma/client';
import { rpc, scValToNative, xdr } from '@stellar/stellar-sdk';
Expand Down Expand Up @@ -126,6 +126,36 @@ const webhookWorker = createWebhookWorker('indexer-webhooks', connectionParams,
});
const getActiveWebhookJob = trackActiveJob(webhookWorker);

// ── Dead-letter queue (DLQ) for webhooks that exhaust all retries (#354) ─────
const DLQ_QUEUE_NAME = 'indexer-webhooks-dlq';
const dlqQueue = new Queue(DLQ_QUEUE_NAME, {
connection: connectionParams,
defaultJobOptions: {
removeOnComplete: { count: 100 },
removeOnFail: { count: 1000 },
},
});

webhookWorker.on('failed', async (job, err) => {
if (!job || job.attemptsMade < (job.opts.attempts ?? WEBHOOK_DEFAULTS.attempts)) return;
// All retries exhausted — move to DLQ
try {
await dlqQueue.add('failed-delivery', {
...job.data,
failedAt: new Date().toISOString(),
error: err?.message ?? String(err),
attempts: job.attemptsMade,
originalJobId: job.id,
} as any);
fastify.log.warn(
{ jobId: job.id, url: job.data.url },
'[Indexer] Webhook moved to dead-letter queue after all retries',
);
} catch (dlqErr) {
fastify.log.error({ err: dlqErr }, '[Indexer] Failed to enqueue job to DLQ');
}
});

// #386 — exponential backoff retry strategy
const redisHealth = createRedisClient(env.REDIS_URL, fastify.log);
redisHealth.on('error', (err) => fastify.log.warn({ err: err.message }, '[Indexer] Redis health client error'));
Expand Down Expand Up @@ -236,20 +266,31 @@ const replayWorker = new Worker(
if (existing) continue;
}

await prisma.indexedEvent.create({
data: {
id: 'evt_' + crypto.randomUUID().replace(/-/g, ''),
stellarId,
contractId: resolvedContractId,
contractName,
topics,
type: topics[0],
rawValue,
decodedPayload: decodedPayload !== null ? (decodedPayload as any) : undefined,
ledger: evt.ledger,
indexedAt: new Date(),
},
});
try {
await prisma.indexedEvent.create({
data: {
id: 'evt_' + crypto.randomUUID().replace(/-/g, ''),
stellarId,
contractId: resolvedContractId,
contractName,
topics,
type: topics[0],
rawValue,
decodedPayload: decodedPayload !== null ? (decodedPayload as any) : undefined,
ledger: evt.ledger,
indexedAt: new Date(),
},
});
} catch (err: any) {
if (err?.code === 'P2002') {
fastify.log.debug(
{ stellarId, contractId: resolvedContractId, ledger: evt.ledger },
'[Indexer] Replay duplicate event — skipping (composite constraint)',
);
continue;
}
throw err;
}

processedLedgers = Math.max(processedLedgers, evt.ledger - fromLedger + 1);
await updateReplayProgress(job.id!, {
Expand Down Expand Up @@ -365,23 +406,35 @@ export async function persistEvent(
rawValue: string,
decodedPayload: unknown,
ledger: number
): Promise<Record<string, unknown>> {
): Promise<Record<string, unknown> | null> {
const id = 'evt_' + crypto.randomUUID().replace(/-/g, '');

const record = await prisma.indexedEvent.create({
data: {
id,
stellarId,
contractId,
contractName,
topics,
type,
rawValue,
decodedPayload: decodedPayload !== null ? (decodedPayload as any) : undefined,
ledger,
indexedAt: new Date(),
},
});
let record: Record<string, unknown>;
try {
record = await prisma.indexedEvent.create({
data: {
id,
stellarId,
contractId,
contractName,
topics,
type,
rawValue,
decodedPayload: decodedPayload !== null ? (decodedPayload as any) : undefined,
ledger,
indexedAt: new Date(),
},
}) as Record<string, unknown>;
} catch (err: any) {
if (err?.code === 'P2002') {
fastify.log.debug(
{ stellarId, contractId, ledger, type },
'[Indexer] Duplicate event — skipping (composite constraint)',
);
return null;
}
throw err;
}

fastify.log.info({ id, type, contractName, ledger }, '[Indexer] Event indexed');

Expand All @@ -396,7 +449,11 @@ export async function persistEvent(

const subs = cacheState.subscriptions.data;
for (const sub of subs) {
await webhookQueue.add('deliver', { url: sub.url, event: record as Record<string, unknown> });
await webhookQueue.add('deliver', {
url: sub.url,
event: record as Record<string, unknown>,
signingSecret: sub.signingSecret ?? undefined,
});
}

return record as Record<string, unknown>;
Expand Down Expand Up @@ -556,6 +613,58 @@ fastify.delete<{ Params: { id: string } }>('/api/webhooks/:id', { preValidation:
return reply.code(204).send();
});

// ── Admin: dead-letter queue (#354) ──────────────────────────────────────────

fastify.get('/api/admin/webhooks/dead-letter', { preValidation: [fastify.serviceAuth] }, async (request) => {
const { limit = 50, offset = 0 } = (request.query as Record<string, unknown>) as { limit?: number; offset?: number };
const safeLimit = Math.min(Number(limit) || 50, 200);
const safeOffset = Math.max(Number(offset) || 0, 0);

const jobs = await dlqQueue.getJobs(['failed'], safeOffset, safeOffset + safeLimit - 1);
const total = await dlqQueue.getJobCounts('failed');

return {
data: jobs.map((job) => ({
id: job.id,
url: job.data.url,
event: job.data.event,
failedAt: (job.data as any).failedAt,
error: (job.data as any).error,
attempts: (job.data as any).attempts,
originalJobId: (job.data as any).originalJobId,
timestamp: job.timestamp,
})),
pagination: { total: total.failed, limit: safeLimit, offset: safeOffset },
};
});

fastify.post<{ Params: { id: string } }>(
'/api/admin/webhooks/dead-letter/:id/replay',
{ preValidation: [fastify.serviceAuth] },
async (request, reply) => {
const { id } = request.params;
const job = await dlqQueue.getJob(id);
if (!job) {
return reply.code(404).send({
error: { code: 'NOT_FOUND', message: `DLQ job ${id} not found` },
});
}

// Re-enqueue on the main webhook delivery queue
await webhookQueue.add('deliver', {
url: job.data.url,
event: job.data.event,
signingSecret: job.data.signingSecret,
});

// Remove from DLQ
await job.remove();

fastify.log.info({ jobId: id, url: job.data.url }, '[Indexer] DLQ job replayed');
return { status: 'requeued', jobId: id, url: job.data.url };
},
);

// ── Stellar RPC polling loop ──────────────────────────────────────────────────

async function pollEvents() {
Expand Down Expand Up @@ -598,8 +707,8 @@ async function pollEvents() {
const contractName = getContractName(resolvedContractId);
const stellarId = typeof evt.id === 'string' ? evt.id : null;

await persistEvent(stellarId, topics, topics[0], resolvedContractId, contractName, rawValue, decodedPayload, evt.ledger);
if (latestLedgerCursor !== undefined) {
const result = await persistEvent(stellarId, topics, topics[0], resolvedContractId, contractName, rawValue, decodedPayload, evt.ledger);
if (latestLedgerCursor !== undefined && result !== null) {
latestLedgerCursor = Math.max(latestLedgerCursor, evt.ledger + 1);
}
}
Expand Down Expand Up @@ -692,6 +801,60 @@ export function stopCleanupScheduler(): void {
}
// ── Startup ───────────────────────────────────────────────────────────────────

/**
* Discovers the correct starting ledger for the polling loop (#352).
*
* Priority:
* 1. INDEX_FROM_LEDGER env var (manual override)
* 2. Latest indexed event ledger + 1 (resume where we left off)
* 3. Network tip - INITIAL_BACKFILL_LEDGERS (fresh deployment)
* 4. Ledger 1 (RPC failure fallback)
*/
export async function discoverStartLedger(): Promise<number> {
// 1. Manual override
if (env.INDEX_FROM_LEDGER) {
const manual = parseInt(env.INDEX_FROM_LEDGER, 10);
if (Number.isFinite(manual) && manual >= 1) {
fastify.log.info({ ledger: manual }, '[Indexer] Starting from manual INDEX_FROM_LEDGER');
return manual;
}
fastify.log.warn({ raw: env.INDEX_FROM_LEDGER }, '[Indexer] Invalid INDEX_FROM_LEDGER — ignoring');
}

// 2. Resume from latest indexed event
try {
const latest = await prisma.indexedEvent.findFirst({
orderBy: { ledger: 'desc' },
select: { ledger: true },
});
if (latest) {
const resumeFrom = latest.ledger + 1;
fastify.log.info({ ledger: resumeFrom, latestIndexed: latest.ledger }, '[Indexer] Resuming from latest indexed event');
return resumeFrom;
}
} catch (err) {
fastify.log.warn({ err: String(err) }, '[Indexer] Failed to query latest indexed event');
}

// 3. Fresh deployment — start from network tip minus backfill window
try {
const tip = await server.getLatestLedger();
const backfill = env.INITIAL_BACKFILL_LEDGERS;
const startLedger = Math.max(1, tip.sequence - backfill);
fastify.log.info(
{ tip: tip.sequence, backfill, startLedger },
'[Indexer] Fresh deployment — starting from network tip minus backfill',
);
return startLedger;
} catch (err) {
fastify.log.warn({ err: String(err) }, '[Indexer] Failed to query Stellar RPC for tip — falling back to ledger 1');
}

// 4. Fallback
fastify.log.warn('[Indexer] No indexed events and RPC unavailable — starting from ledger 1');
return 1;
}

const start = async () => {
try {
// #391 — wait for both dependencies before accepting traffic
Expand All @@ -701,6 +864,9 @@ const start = async () => {
// #387 — Redis memory monitoring
startRedisMemoryMonitor(redisHealth, fastify.log);

// #352 — smart startup ledger discovery
latestLedgerCursor = await discoverStartLedger();

await fastify.listen({ port: PORT, host: '0.0.0.0' });
fastify.log.info('[Indexer] Starting Stellar RPC polling loop...');
pollEvents();
Expand All @@ -715,9 +881,10 @@ process.on('SIGTERM', async () => {
await prisma.$disconnect();
await replayQueue.close();
await closeWorkerWithTimeout(replayWorker, 'indexer-replays', fastify.log, getActiveReplayJob);
await replayProgressRedis.quit().catch(() => {});
await webhookQueue.close();
await closeWorkerWithTimeout(webhookWorker, 'indexer-webhooks', fastify.log, getActiveWebhookJob);
await dlqQueue.close();
await replayProgressRedis.quit().catch(() => {});
await fastify.close();
await new Promise<void>((resolve) => metricsServer.close(() => resolve()));
process.exit(0);
Expand Down
Loading
Loading