diff --git a/.env.example b/.env.example index b7d00a5..1cec8b9 100644 --- a/.env.example +++ b/.env.example @@ -4,9 +4,20 @@ KEEPER_SECRET= # Optional public key. If set, the keeper will verify it matches KEEPER_SECRET. KEEPER_PUBLIC_KEY= -# Network passphrase. Defaults to the Stellar testnet. +# Network passphrase. Defaults to the Stellar mainnet. # For mainnet, set to: "Public Global Stellar Network ; September 2015" NETWORK_PASSTHRASE= # Must be "true" when using the mainnet passphrase. ALLOW_MAINNET=false + +# Indexer backfill mode: optional start and end ledger for controlled backfill. +# If unset, the indexer resumes from the last checkpoint. +BACKFILL_START_LEDGER= +BACKFILL_END_LEDGER= + +# Backfill rate limit (requests per second). Defaults to a safe value compatible with RPC. +BACKFILL_RATE_LIMIT= + +# Path to the checkpoint file for progressive backfill checkpoints. +BACKFILL_CHECKPOINT_FILE= diff --git a/config.ts b/config.ts index fe59bf5..7d12fd5 100644 --- a/config.ts +++ b/config.ts @@ -18,7 +18,7 @@ export interface Config { /** * Loads configuration from environment variables with defaults. */ -export function loadConfig(env: NodeJS.ProcessEnv = process.env): Config { +export function loadAlertConfig(env: NodeJS.ProcessEnv = process.env): Config { return { webhookUrl: env.WEBHOOK_URL || '', dedupWindowMs: Number(env.DEDUP_WINDOW_MS || 3600000), @@ -84,3 +84,26 @@ export function formatConfig(config: KeeperConfig): string { const redacted = { ...config, secret: redactSecret(config.secret) }; return JSON.stringify(redacted, null, 2); } + +export interface BackfillConfig { + /** Starting ledger sequence number for backfill. 0 means from the last checkpoint. */ + startLedger: number; + /** Ending ledger sequence number for backfill (inclusive). 0 means to the latest. */ + endLedger: number; + /** Path to the checkpoint file for progressive checkpointing. */ + checkpointFile: string; + /** Maximum number of ledger fetch requests per second. 0 means no limit. */ + rateLimit: number; +} + +/** + * Loads backfill configuration from environment variables with defaults. + */ +export function loadBackfillConfig(env: NodeJS.ProcessEnv = process.env): BackfillConfig { + return { + startLedger: Number(env.BACKFILL_START_LEDGER || 0), + endLedger: Number(env.BACKFILL_END_LEDGER || 0), + checkpointFile: env.BACKFILL_CHECKPOINT_FILE || '', + rateLimit: Number(env.BACKFILL_RATE_LIMIT || 0), + }; +} diff --git a/indexer.ts b/indexer.ts new file mode 100644 index 0000000..e45ebbc --- /dev/null +++ b/indexer.ts @@ -0,0 +1 @@ +export async function backfillLedgers(o:any,s:number,e:number){const c=o.checkpointEvery??10;const r=o.rateLimit?.requestsPerSecond??5;let l=await o.checkpointStore.getLastCheckpoint();let cur=s;if(l!=null&&l>=s-1)cur=Math.max(s,l+1);if(cur>e)return;let f=0,u=0;for(;cur<=e;cur++){await new Promise(z=>setTimeout(z,1000/Math.max(1,r)));const ev=await o.rpc.getEvents(cur);u+=await o.eventStore.upsertEvents(ev);f++;if(cur%c===0||cur===e)await o.checkpointStore.setLastCheckpoint(cur);}await o.checkpointStore.setLastCheckpoint(e);}export function createRpcClient(url:string,h:Record{const x=await fetch(url,{method:"POST",headers:{"content-type":"application/json",...h},body:JSON.stringify({method:m,params:p})});if(!x.ok)throw new Error(m+" "+x.status);const j=await x.json();if(j.error)throw new Error(j.error);return j.result;};return {getLedger:n=>call("getLedger",{sequence:n}),getEvents:n=>call("getEvents",{ledgerSequence:n})}} \ No newline at end of file diff --git a/query-events.ts b/query-events.ts new file mode 100644 index 0000000..386de2d --- /dev/null +++ b/query-events.ts @@ -0,0 +1,48 @@ +import { promises as fs } from 'fs'; +import * as path from 'path'; +import { fetchLedgerEvents } from './indexer'; + +export interface Options { + startLedger?: number; + endLedger: number; + checkpointFile?: string; + dataFile?: string; + rateLimitPerSecond?: number; +} + +export async function backfillEvents(opts: Options) { + const cp = path.resolve(opts.checkpointFile ?? '.checkpoint.json'); + const data = path.resolve(opts.dataFile ?? '.events.json'); + const checkpoint = await readJson(cp); + const start = opts.startLedger ?? (checkpoint?.last ?? 0) + 1; + const end = opts.endLedger; + const store = (await readJson(data)) ?? {}; + const delay = 1000 / (opts.rateLimitPerSecond ?? 10); + let upserts = 0; + + for (let ledger = start; ledger <= end; ledger++) { + await sleep(delay); + const events = await fetchLedgerEvents(ledger); + for (const event of events) { + const id = String(event.id ?? `${ledger}:${event.sequence ?? 0}`); + if (!store[id]) { + store[id] = event; + upserts++; + } + } + await writeJson(cp, { last: ledger, updatedAt: new Date().toISOString() }); + } + + await writeJson(data, store); + return { processedLedgers: end - start + 1, upsertedEvents: upserts }; +} + +function sleep(ms: number) { return new Promise((r) => setTimeout(r, ms)); } + +async function readJson(file: string): Promise { + try { return JSON.parse(await fs.readFile(file, 'utf8')); } catch { return undefined; } +} + +async function writeJson(file: string, value: any): Promise { + await fs.writeFile(file, JSON.stringify(value, null, 2)); +} diff --git a/replay-events.ts b/replay-events.ts new file mode 100644 index 0000000..328122a --- /dev/null +++ b/replay-events.ts @@ -0,0 +1,64 @@ +import { readFile, writeFile, rename, mkdir } from 'fs/promises'; +import path from 'path'; +import { Pool } from 'pg'; +import { setTimeout as sleep } from 'timers/promises'; + +const RPC_URL = process.env.RPC_URL!; +const DATABASE_URL_= process.env.DATABASE_URL; +const START = Number(process.env.START_LEDGER!); +const END = Number(process.env.END_LEDGER!); +const CHECKPOINT = process.env.CHECKPOINT_FILE || path.join(process.cwd(), '.replay-checkpoint.json'); +const RATE_LIMIT_MS = Number(process.env.RATE_LIMIT_MS ?? 100); + +async function getEvents(rpcUrl: string, ledger: number): Promise { + const res = await fetch(rpcUrl, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'getEvents', params: [{ startLedger: ledger, endLedger: ledger }] }), + }); + const data: any = await res.json(); + if (data.error) throw new Error(data.error.message); + let events = data.result?.events ?? []; + let cursor = data.result?.cursor; + while (cursor) { + const res2 = await fetch(rpcUrl, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'getEvents', params: [{ startLedger: ledger, endLedger: ledger, cursor }] }), + }); + const data2: any = await res2.json(); + events = events.concat(data2.result?.events ?? []); + cursor = data2.result?.cursor; + } + return events; +} + +async function main() { + const pool = new Pool({ connectionString: DATABASE_URL }); + let last = null; + try { + const raw = await readFile(CHECKPOINT, 'utf-8'); + last = JSON.parse(raw).last; + } catch {} + const start = last === null ? START : last + 1; + if (start > END) { console.log('Already backfilled'); await pool.end(); return; } + + for (let ledger = start; ledger <= END; ledger++) { + await sleep(RATE_LIMIT_MS); + const events = await getEvents(RPC_URL, ledger); + for (const e of events) { + // Idempotent upsert: assume primary key is id + await pool.query( + 'INSERT INTO events (id, ledger, data) VALUES ($1, $2, $3) ON CONFLICT (id) DO NOTHING', + [e.id, ledger, JSON.stringify(e)] + ); + } + await mkdiv(path.dirname(CHECKPOINT), { recursive: true }); + await writeFile(CHECKPOINT + '.tmp', JSON.stringify({ last: ledger })); + await rename(CHECKPOINT + '.tmp', CHECKPOINT); + console.log(`Processed ledger ${ledger}, events=${events.length}); + } + await pool.end(); +} + +main().catch(err => { console.error(err); process.exit(1); });