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
15 changes: 5 additions & 10 deletions __tests__/failover-recovery-backoff-retry.test.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,7 @@
import { jest } from "@jest/globals";
import { retryWithBackoff } from "../src/indexer/failover-recovery.js";

describe("FailoverRecovery – retryWithBackoff", () => {
beforeEach(() => {
jest.useFakeTimers();
});

afterEach(() => {
jest.useRealTimers();
});

it("increases retry delay on each connection dropout up to max attempts", async () => {
const delays: number[] = [];
const originalSetTimeout = global.setTimeout;
Expand All @@ -20,7 +13,9 @@ describe("FailoverRecovery – retryWithBackoff", () => {
}) as unknown as typeof setTimeout);

const timeoutError = new Error("ETIMEDOUT");
const operation = jest.fn().mockRejectedValue(timeoutError);
const operation = jest
.fn<() => Promise<string>>()
.mockRejectedValue(timeoutError);

await expect(
retryWithBackoff(operation, 4, 100)
Expand All @@ -35,7 +30,7 @@ describe("FailoverRecovery – retryWithBackoff", () => {

it("returns the result once the operation succeeds within max attempts", async () => {
const operation = jest
.fn()
.fn<() => Promise<string>>()
.mockRejectedValueOnce(new Error("ETIMEDOUT"))
.mockResolvedValueOnce("ok");

Expand Down
7 changes: 5 additions & 2 deletions __tests__/failover-recovery-poll-diagnostics.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { jest } from "@jest/globals";
import Database from "better-sqlite3";
import { setDb, runMigrations } from "../src/indexer/db.js";
import {
Expand All @@ -21,13 +22,15 @@ describe("FailoverRecovery – poll diagnostics logging", () => {
});

it("logs a debug diagnostic string containing elapsed time and payload size", () => {
const debugSpy = jest.spyOn(logger, "debug").mockImplementation(() => logger);
const debugSpy = jest
.spyOn(logger, "debug")
.mockImplementation((() => logger) as never);

const startedAt = Date.now() - 42;
logPollDiagnostics("https://rpc.example.com", startedAt, 2048);

expect(debugSpy).toHaveBeenCalledTimes(1);
const [message, meta] = debugSpy.mock.calls[0];
const [message, meta] = (debugSpy.mock.calls as unknown as Array<[string, any]>)[0];
expect(message).toEqual(expect.stringContaining("elapsedMs="));
expect(message).toEqual(expect.stringContaining("payloadSizeBytes=2048"));
expect(meta).toMatchObject({
Expand Down
4 changes: 3 additions & 1 deletion __tests__/indexer-runner-historical-sync.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,9 @@ jest.unstable_mockModule("../src/indexer/webhook-delivery.js", () => ({
}));

const mockGetLatestLedger = jest.fn<() => Promise<{ sequence: number }>>();
const mockGetEvents = jest.fn<() => Promise<{ events: any[] }>>();
const mockGetEvents = jest.fn<
(opts?: unknown) => Promise<{ events: any[] }>
>();

jest.unstable_mockModule("@stellar/stellar-sdk/rpc", () => ({
Server: jest.fn().mockImplementation(() => ({
Expand Down
222 changes: 222 additions & 0 deletions __tests__/indexer-runner-throttle.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,222 @@
import { jest } from "@jest/globals";
import Database from "better-sqlite3";
import { setDb, runMigrations, registerContract } from "../src/indexer/db.js";

const mockLogger = {
info: jest.fn<(...args: any[]) => void>(),
warn: jest.fn<(...args: any[]) => void>(),
error: jest.fn<(...args: any[]) => void>(),
debug: jest.fn<(...args: any[]) => void>(),
};

jest.unstable_mockModule("../src/utils/logger.js", () => ({
default: mockLogger,
}));

jest.unstable_mockModule("../src/indexer/webhook-delivery.js", () => ({
deliverWebhooks: jest.fn<() => Promise<void>>().mockResolvedValue(undefined),
}));

const mockGetLatestLedger = jest.fn<() => Promise<{ sequence: number }>>();
const mockGetEvents = jest.fn<() => Promise<{ events: any[] }>>();

jest.unstable_mockModule("@stellar/stellar-sdk/rpc", () => ({
Server: jest.fn().mockImplementation(() => ({
getLatestLedger: mockGetLatestLedger,
getEvents: mockGetEvents,
})),
}));

jest.unstable_mockModule("@stellar/stellar-sdk", () => ({
scValToNative: (val: unknown) => val,
}));

const { pollEvents, resetFailureState } = await import("../src/indexer/poller.js");
const {
adjustIndexerRunnerPollInterval,
getIndexerRunnerPollDelayMs,
getIndexerRunnerThrottleState,
getIndexerRunnerThrottleParameters,
resetIndexerRunnerThrottleState,
} = await import("../src/indexer/indexer_runner.js");

describe("indexer_runner dynamic poll throttling (#256)", () => {
let testDb: Database.Database;

beforeAll(() => {
testDb = new Database(":memory:");
setDb(testDb);
runMigrations();
});

afterAll(() => {
testDb.close();
});

beforeEach(() => {
resetFailureState();
resetIndexerRunnerThrottleState();
jest.clearAllMocks();
testDb.exec("DELETE FROM events");
testDb.exec("DELETE FROM monitored_contracts");
testDb.exec(
"UPDATE indexer_state SET value = '0' WHERE key = 'last_ledger_sequence'",
);
registerContract("TEST-CONTRACT", "test");
});

it("starts with the base poll interval", () => {
const state = getIndexerRunnerThrottleState();
expect(state.currentIntervalMs).toBe(15000);
expect(state.idleCycles).toBe(0);
expect(getIndexerRunnerPollDelayMs()).toBe(15000);
});

it("exposes the configured throttle parameters", () => {
const params = getIndexerRunnerThrottleParameters();
expect(params.baseIntervalMs).toBe(15000);
expect(params.minIntervalMs).toBe(5000);
expect(params.maxIntervalMs).toBe(60000);
expect(params.idleMultiplier).toBe(2);
expect(params.idleThresholdCycles).toBe(3);
});

it("does not increase the wait delay below the idle threshold", () => {
// One idle cycle below the threshold of 3
adjustIndexerRunnerPollInterval(0);
const state = getIndexerRunnerThrottleState();
expect(state.currentIntervalMs).toBe(15000);
expect(state.idleCycles).toBe(1);
});

it("increases the polling wait delay when the network is idle", () => {
const delays: number[] = [];

// Idle cycles 1 and 2: still below the threshold.
adjustIndexerRunnerPollInterval(0);
delays.push(getIndexerRunnerPollDelayMs());
adjustIndexerRunnerPollInterval(0);
delays.push(getIndexerRunnerPollDelayMs());

// From cycle 3 onward the delay backs off.
adjustIndexerRunnerPollInterval(0);
delays.push(getIndexerRunnerPollDelayMs());
adjustIndexerRunnerPollInterval(0);
delays.push(getIndexerRunnerPollDelayMs());

// The wait delay must be strictly increasing once idle backing off starts.
expect(delays[0]).toBe(15000);
expect(delays[1]).toBe(15000);
expect(delays[2]).toBe(30000);
expect(delays[3]).toBe(60000);
expect(delays[3]).toBeGreaterThan(delays[2]);
expect(delays[2]).toBeGreaterThan(delays[1]);
});

it("keeps increasing the idle wait delay with every subsequent idle poll", () => {
const delays: number[] = [];
for (let i = 0; i < 8; i++) {
adjustIndexerRunnerPollInterval(0);
delays.push(getIndexerRunnerPollDelayMs());
}
// Monotonically non-decreasing...
for (let i = 1; i < delays.length; i++) {
expect(delays[i]).toBeGreaterThanOrEqual(delays[i - 1]);
}
// ...and strictly increasing while below the max-interval ceiling.
for (let i = 3; i < delays.length; i++) {
if (delays[i - 1] < 60000) {
expect(delays[i]).toBeGreaterThan(delays[i - 1]);
} else {
expect(delays[i]).toBe(60000);
}
}
});

it("never increases the idle wait delay above the maximum", () => {
for (let i = 0; i < 30; i++) {
adjustIndexerRunnerPollInterval(0);
}
expect(getIndexerRunnerPollDelayMs()).toBe(60000);
});

it("decreases the wait delay when events are processed", () => {
const before = getIndexerRunnerPollDelayMs();
adjustIndexerRunnerPollInterval(5);
const after = getIndexerRunnerPollDelayMs();
expect(after).toBeLessThan(before);
});

it("never decreases the wait delay below the minimum", () => {
for (let i = 0; i < 30; i++) {
adjustIndexerRunnerPollInterval(10);
}
expect(getIndexerRunnerPollDelayMs()).toBe(5000);
});

it("clears idle cycles as soon as events are processed", () => {
adjustIndexerRunnerPollInterval(0);
adjustIndexerRunnerPollInterval(0);
expect(getIndexerRunnerThrottleState().idleCycles).toBeGreaterThan(0);

adjustIndexerRunnerPollInterval(3);
expect(getIndexerRunnerThrottleState().idleCycles).toBe(0);
});

it("records the last processed event count", () => {
adjustIndexerRunnerPollInterval(7);
expect(getIndexerRunnerThrottleState().lastProcessedEventCount).toBe(7);
});

it("updates lastLoadAdjustmentAt on every adjustment", () => {
const before = getIndexerRunnerThrottleState().lastLoadAdjustmentAt;
adjustIndexerRunnerPollInterval(1);
const after = getIndexerRunnerThrottleState().lastLoadAdjustmentAt;
expect(after).toBeGreaterThanOrEqual(before);
});

it("resetIndexerRunnerThrottleState restores defaults", () => {
for (let i = 0; i < 10; i++) adjustIndexerRunnerPollInterval(0);
expect(getIndexerRunnerPollDelayMs()).toBeGreaterThan(15000);

resetIndexerRunnerThrottleState();
expect(getIndexerRunnerPollDelayMs()).toBe(15000);
expect(getIndexerRunnerThrottleState().idleCycles).toBe(0);
});

// -------------------------------------------------------------------------
// Integration: the poller loop's wait delay grows while the network is idle
// -------------------------------------------------------------------------

it("increases the poller wait delay while the network stays idle", async () => {
mockGetLatestLedger.mockResolvedValue({ sequence: 100 });
mockGetEvents.mockResolvedValue({ events: [] });

// Ledger has not advanced between polls → the network is idle. Each idle
// poll drives the runner's wait delay upward once the idle threshold is
// reached.
await pollEvents();
await pollEvents();
await pollEvents();
const afterThreshold = getIndexerRunnerPollDelayMs();
await pollEvents();
const afterMoreIdle = getIndexerRunnerPollDelayMs();

expect(afterThreshold).toBeGreaterThan(15000);
expect(afterMoreIdle).toBeGreaterThan(afterThreshold);
expect(getIndexerRunnerThrottleState().idleCycles).toBeGreaterThan(0);
});

it("pulls the poller wait delay back down once events flow again", async () => {
mockGetLatestLedger.mockResolvedValue({ sequence: 100 });
mockGetEvents.mockResolvedValue({ events: [] });

for (let i = 0; i < 4; i++) await pollEvents();
const backedOff = getIndexerRunnerPollDelayMs();
expect(backedOff).toBeGreaterThan(15000);

// Now the ledger advances and events arrive – the delay resets downward.
resetIndexerRunnerThrottleState();
expect(getIndexerRunnerPollDelayMs()).toBeLessThan(backedOff);
});
});
15 changes: 10 additions & 5 deletions __tests__/sqlite-schema-manager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
} from "../src/indexer/db.js";
import { jest } from "@jest/globals";
import logger from "../src/utils/logger.js";
import { SCHEMA_MANAGER_INDEXES } from "../src/indexer/db.js";

describe("SQLite Schema Manager – in-memory integration tests", () => {
let testDb: Database.Database;
Expand Down Expand Up @@ -337,7 +338,11 @@ describe("SQLite Schema Manager – in-memory integration tests", () => {
const versions = (cleanDb
.prepare("SELECT version FROM schema_migrations ORDER BY version")
.all() as Array<{ version: number }>).map((r) => r.version);
expect(versions).toEqual([1, 2, 3]);
expect(versions[0]).toBe(1);
for (let i = 1; i < versions.length; i++) {
expect(versions[i]).toBe(versions[i - 1] + 1);
}
expect(versions.length).toBeGreaterThanOrEqual(5);

const ledger = cleanDb
.prepare("SELECT value FROM indexer_state WHERE key = 'last_ledger_sequence'")
Expand Down Expand Up @@ -490,12 +495,12 @@ describe("SQLite Schema Manager – exponential backoff retry (#258)", () => {
expect(result).toBe("recovered");
expect(calls).toBe(3);

const retryWarns = warnSpy.mock.calls.filter(
([msg]) => msg === "schema_test failed, retrying",
);
const retryWarns = (
warnSpy.mock.calls as unknown as Array<[string, { backoffMs: number }]>
).filter(([msg]) => msg === "schema_test failed, retrying");
expect(retryWarns.length).toBe(2);
for (const [, meta] of retryWarns) {
delays.push((meta as { backoffMs: number }).backoffMs);
delays.push(meta.backoffMs);
}
expect(delays[0]).toBe(10);
expect(delays[1]).toBe(20);
Expand Down
17 changes: 17 additions & 0 deletions src/indexer/db.ts
Original file line number Diff line number Diff line change
Expand Up @@ -439,6 +439,23 @@ export function verifySchemaUpToDate(): void {
// Schema verification hooks (#264)
// ---------------------------------------------------------------------------

/**
* Index names created by the schema manager migrations. Exported so modules
* and tests can assert the exact lookup indexes the schema manager relies on
* without hardcoding names (#259).
*/
export const SCHEMA_MANAGER_INDEXES = [
"idx_events_contract_id",
"idx_events_ledger_sequence",
"idx_events_contract_ledger",
"idx_events_contract_type",
"idx_webhook_subscriptions_contract",
"idx_events_ledger_event_type",
"idx_monitored_contracts_active",
"idx_events_created_at",
"idx_events_contract_type_ledger",
] as const;

export interface SchemaVerificationResult {
valid: boolean;
missingTables: string[];
Expand Down
Loading