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
21 changes: 14 additions & 7 deletions __tests__/indexer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -60,15 +60,22 @@ describe("Indexer Database", () => {
});

it("does not re-apply already-applied migrations (idempotent)", () => {
const before = testDb
.prepare("SELECT version FROM schema_migrations ORDER BY version")
.all() as Array<{ version: number }>;

// Running again should not throw and should not duplicate rows
runMigrations();
const rows = testDb
.prepare("SELECT version FROM schema_migrations")
.all();
// We ship 5 migrations (events/indexer_state + monitored_contracts + indexes +
// ledger range indexes + schema-manager lookup indexes)
const versions = [...new Set((rows as any[]).map((r) => r.version))];
expect(versions.length).toBe(5);

const after = testDb
.prepare("SELECT version FROM schema_migrations ORDER BY version")
.all() as Array<{ version: number }>;
expect(after).toEqual(before);

// The full migration set is applied exactly once (sequential from 1).
const versions = after.map((r) => r.version);
expect(new Set(versions).size).toBe(versions.length);
expect(Math.min(...versions)).toBe(1);
});
});

Expand Down
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
146 changes: 146 additions & 0 deletions __tests__/sqlite-vacuum-indexes.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
import Database from "better-sqlite3";
import { runMigrations, setDb } from "../src/indexer/db.js";
import {
VACUUM_CLEANER_INDEXES,
getVacuumIndexNames,
ensureVacuumIndexes,
vacuumExplainQueryPlan,
vacuumQueryPlanUsesIndex,
} from "../src/indexer/sqlite_vacuum_cleaner.js";

describe("sqlite_vacuum_cleaner – SQLite index structures (#344)", () => {
let db: Database.Database;

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

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

beforeEach(() => {
db.exec("DELETE FROM events");
seedEvents();
});

/** Seed enough rows that the planner prefers an index scan over a full scan. */
function seedEvents(): void {
const insert = db.prepare(
`INSERT OR IGNORE INTO events
(contract_id, event_type, ledger_sequence, timestamp, data_json, created_at)
VALUES (?, 'funded', ?, ?, '{}', ?)`,
);
const now = Date.now();
for (let i = 0; i < 300; i++) {
const createdDaysAgo = (i % 300);
insert.run(
`C${i % 10}`,
1000 + i,
1_700_000_000 + i,
new Date(now - createdDaysAgo * 86_400_000).toISOString(),
);
}
}

it("migration 6 creates all vacuum cleaner lookup indexes", () => {
const rows = db
.prepare(
`SELECT name FROM sqlite_master
WHERE type = 'index' AND name IN (${getVacuumIndexNames()
.map(() => "?")
.join(", ")})`,
)
.all(...getVacuumIndexNames()) as Array<{ name: string }>;
const names = rows.map((r) => r.name);

for (const indexName of getVacuumIndexNames()) {
expect(names).toContain(indexName);
}
});

it("ensureVacuumIndexes is idempotent and returns every managed index", () => {
const first = ensureVacuumIndexes(db);
expect(first).toEqual(getVacuumIndexNames());

const second = ensureVacuumIndexes(db);
expect(second).toEqual(getVacuumIndexNames());
expect(second).toEqual(expect.arrayContaining(first));
});

it("retention-time lookups (created_at) use the vacuum cleaner index", () => {
// Mirror pruneOldEvents' predicate (DELETE ... WHERE created_at < now-N).
const deletePlan = vacuumExplainQueryPlan(
db,
`DELETE FROM events WHERE created_at < datetime('now', ?)`,
"-90 days",
);
expect(
vacuumQueryPlanUsesIndex(deletePlan, VACUUM_CLEANER_INDEXES.eventsCreatedAt),
).toBe(true);

// SELECT-form of the same lookup also resolves through a managed index.
const selectPlan = vacuumExplainQueryPlan(
db,
`SELECT * FROM events WHERE created_at < datetime('now', ?)`,
"-90 days",
);
const used = getVacuumIndexNames().some((name) =>
vacuumQueryPlanUsesIndex(selectPlan, name),
);
expect(used).toBe(true);
});

it("ledger-range lookups use the vacuum cleaner index", () => {
// Mirror pruneEventsInLedgerRange's predicate (ledger_sequence BETWEEN).
const deletePlan = vacuumExplainQueryPlan(
db,
`DELETE FROM events WHERE ledger_sequence >= ? AND ledger_sequence <= ?`,
1005,
1020,
);
expect(
vacuumQueryPlanUsesIndex(
deletePlan,
VACUUM_CLEANER_INDEXES.eventsLedgerSequence,
),
).toBe(true);

const selectPlan = vacuumExplainQueryPlan(
db,
`SELECT * FROM events WHERE ledger_sequence >= ? AND ledger_sequence <= ?`,
1005,
1020,
);
const used = getVacuumIndexNames().some((name) =>
vacuumQueryPlanUsesIndex(selectPlan, name),
);
expect(used).toBe(true);
});

it("asserts the combined retention + range pruning lookup uses a managed index", () => {
const plan = vacuumExplainQueryPlan(
db,
`SELECT ledger_sequence FROM events
WHERE created_at < datetime('now', ?) AND ledger_sequence >= ?`,
"-30 days",
1000,
);
const used = getVacuumIndexNames().some((name) =>
vacuumQueryPlanUsesIndex(plan, name),
);
expect(used).toBe(true);
});

it("vacuumQueryPlanUsesIndex is false when no managed index is referenced", () => {
const plan = [{ detail: "SCAN events" }];
expect(vacuumQueryPlanUsesIndex(plan, VACUUM_CLEANER_INDEXES.eventsCreatedAt)).toBe(
false,
);
expect(vacuumQueryPlanUsesIndex(plan, "idx_events_created_at_write_test")).toBe(
false,
);
});
});
34 changes: 34 additions & 0 deletions src/indexer/db.ts
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,23 @@ const MIGRATIONS: Migration[] = [
ON events (contract_id, event_type, ledger_sequence);
`,
},
{
version: 6,
description: "add sqlite_vacuum_cleaner lookup indexes (#344)",
up: `
CREATE INDEX IF NOT EXISTS idx_events_created_at
ON events (created_at);

CREATE INDEX IF NOT EXISTS idx_events_ledger_sequence
ON events (ledger_sequence);

CREATE INDEX IF NOT EXISTS idx_events_created_at_ledger
ON events (created_at, ledger_sequence);

CREATE INDEX IF NOT EXISTS idx_events_ledger_created_at
ON events (ledger_sequence, created_at);
`,
},
];

// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -439,6 +456,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