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
13 changes: 12 additions & 1 deletion .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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=
25 changes: 24 additions & 1 deletion config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down Expand Up @@ -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),
};
}
1 change: 1 addition & 0 deletions indexer.ts
Original file line number Diff line number Diff line change
@@ -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<string,string=={}){const call=async(m:any,p:any)=>{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})}}
48 changes: 48 additions & 0 deletions query-events.ts
Original file line number Diff line number Diff line change
@@ -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<any> {
try { return JSON.parse(await fs.readFile(file, 'utf8')); } catch { return undefined; }
}

async function writeJson(file: string, value: any): Promise<void> {
await fs.writeFile(file, JSON.stringify(value, null, 2));
}
64 changes: 64 additions & 0 deletions replay-events.ts
Original file line number Diff line number Diff line change
@@ -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<any[]> {
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); });
Loading