diff --git a/__tests__/failover-recovery-backoff-retry.test.ts b/__tests__/failover-recovery-backoff-retry.test.ts index 80f8527..c509bd5 100644 --- a/__tests__/failover-recovery-backoff-retry.test.ts +++ b/__tests__/failover-recovery-backoff-retry.test.ts @@ -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; @@ -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>() + .mockRejectedValue(timeoutError); await expect( retryWithBackoff(operation, 4, 100) @@ -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>() .mockRejectedValueOnce(new Error("ETIMEDOUT")) .mockResolvedValueOnce("ok"); diff --git a/__tests__/failover-recovery-poll-diagnostics.test.ts b/__tests__/failover-recovery-poll-diagnostics.test.ts index 47fc14b..6aa97a2 100644 --- a/__tests__/failover-recovery-poll-diagnostics.test.ts +++ b/__tests__/failover-recovery-poll-diagnostics.test.ts @@ -1,3 +1,4 @@ +import { jest } from "@jest/globals"; import Database from "better-sqlite3"; import { setDb, runMigrations } from "../src/indexer/db.js"; import { @@ -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({ diff --git a/__tests__/indexer-runner-historical-sync.test.ts b/__tests__/indexer-runner-historical-sync.test.ts index 8b02ab0..bda27dd 100644 --- a/__tests__/indexer-runner-historical-sync.test.ts +++ b/__tests__/indexer-runner-historical-sync.test.ts @@ -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(() => ({ diff --git a/__tests__/indexer.test.ts b/__tests__/indexer.test.ts index 5f34e5f..fcb21d1 100644 --- a/__tests__/indexer.test.ts +++ b/__tests__/indexer.test.ts @@ -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); }); }); diff --git a/__tests__/sqlite-schema-manager.test.ts b/__tests__/sqlite-schema-manager.test.ts index 1dc7594..b27cbd9 100644 --- a/__tests__/sqlite-schema-manager.test.ts +++ b/__tests__/sqlite-schema-manager.test.ts @@ -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; @@ -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'") @@ -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); diff --git a/__tests__/sqlite-vacuum-indexes.test.ts b/__tests__/sqlite-vacuum-indexes.test.ts new file mode 100644 index 0000000..98dbd8d --- /dev/null +++ b/__tests__/sqlite-vacuum-indexes.test.ts @@ -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, + ); + }); +}); \ No newline at end of file diff --git a/src/indexer/db.ts b/src/indexer/db.ts index 682d497..399421d 100644 --- a/src/indexer/db.ts +++ b/src/indexer/db.ts @@ -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); + `, + }, ]; // --------------------------------------------------------------------------- @@ -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[]; diff --git a/src/indexer/sqlite_vacuum_cleaner.ts b/src/indexer/sqlite_vacuum_cleaner.ts index cb1330e..7d7ea81 100644 --- a/src/indexer/sqlite_vacuum_cleaner.ts +++ b/src/indexer/sqlite_vacuum_cleaner.ts @@ -443,3 +443,87 @@ export function assertVacuumSchemaValid(db: Database.Database): void { ); } } + +// --------------------------------------------------------------------------- +// Issue 5: SQLite index structures for vacuum cleaner lookups (#344) +// --------------------------------------------------------------------------- + +/** + * Index names backing the vacuum cleaner's row lookups. + * + * The cleaner finds stale rows through two predicates: + * 1. retention-pruning by `created_at` → covered by idx_events_created_at + * (and idx_events_created_at_ledger for the composite layout). + * 2. ledger-range pruning by `ledger_sequence` → covered by + * idx_events_ledger_sequence (and idx_events_ledger_created_at). + * + * Exporting the names lets tests run EXPLAIN QUERY PLAN and assert the indexes + * are actually used, and lets operators inspect the schema (#344). + */ +export const VACUUM_CLEANER_INDEXES = { + eventsCreatedAt: "idx_events_created_at", + eventsLedgerSequence: "idx_events_ledger_sequence", + eventsCreatedAtLedger: "idx_events_created_at_ledger", + eventsLedgerCreatedAt: "idx_events_ledger_created_at", +} as const; + +/** All index names managed by the vacuum cleaner. */ +export function getVacuumIndexNames(): string[] { + return Object.values(VACUUM_CLEANER_INDEXES); +} + +/** + * Ensure every index the vacuum cleaner relies on exists, creating any missing + * ones idempotently. Used alongside the schema manager migrations so the cleaner + * can self-heal a database that predates migration 6 (#344). + * + * @returns The names of the indexes that are now present. + */ +export function ensureVacuumIndexes(db: Database.Database): string[] { + db.exec( + ` + CREATE INDEX IF NOT EXISTS ${VACUUM_CLEANER_INDEXES.eventsCreatedAt} + ON events (created_at); + + CREATE INDEX IF NOT EXISTS ${VACUUM_CLEANER_INDEXES.eventsLedgerSequence} + ON events (ledger_sequence); + + CREATE INDEX IF NOT EXISTS ${VACUUM_CLEANER_INDEXES.eventsCreatedAtLedger} + ON events (created_at, ledger_sequence); + + CREATE INDEX IF NOT EXISTS ${VACUUM_CLEANER_INDEXES.eventsLedgerCreatedAt} + ON events (ledger_sequence, created_at); + `, + ); + return getVacuumIndexNames(); +} + +/** + * Run SQLite's EXPLAIN QUERY PLAN for a parameterized statement against the + * given connection. Mirrors ledger_range_tracker's helper so tests can assert + * the vacuum cleaner's indexes are utilized for lookups (#344). + */ +export function vacuumExplainQueryPlan( + db: Database.Database, + sql: string, + ...params: unknown[] +): Array> { + return db.prepare(`EXPLAIN QUERY PLAN ${sql}`).all(...params) as Array< + Record + >; +} + +/** + * True when any EXPLAIN QUERY PLAN detail row references the expected index + * name. Mirrors ledger_range_tracker's helper (#344). + */ +export function vacuumQueryPlanUsesIndex( + plan: Array>, + indexName: string, +): boolean { + return plan.some((row) => + Object.values(row).some( + (value) => typeof value === "string" && value.includes(indexName), + ), + ); +}