diff --git a/__tests__/database-writer-pool-indexes.test.ts b/__tests__/database-writer-pool-indexes.test.ts new file mode 100644 index 0000000..5bae16b --- /dev/null +++ b/__tests__/database-writer-pool-indexes.test.ts @@ -0,0 +1,433 @@ +import Database from "better-sqlite3"; +import { setDb, runMigrations, closeDb, insertEvent } from "../src/indexer/db.js"; +import { + WRITER_POOL_INDEXES, + WRITER_POOL_QUERIES, + WRITER_POOL_UNIQUE_INDEXES, + createSqlOperation, + explainWriterPoolQueryPlan, + queueWrite, + verifyWriterPoolIndexes, + verifyWriterPoolSchema, + writerPoolQueryPlanUsesIndex, + writerPoolQueryPlanUsesTempBTree, +} from "../src/indexer/database-writer-pool.js"; + +describe("database_writer_pool – SQLite index structures (#326)", () => { + let testDb: Database.Database; + + beforeEach(() => { + testDb = new Database(":memory:"); + setDb(testDb); + runMigrations(); + seedWriterPoolLookupRows(testDb); + }); + + afterEach(() => { + closeDb(); + }); + + function indexNames(): string[] { + return ( + testDb + .prepare("SELECT name FROM sqlite_master WHERE type = 'index'") + .all() as Array<{ name: string }> + ).map((row) => row.name); + } + + describe("migrations", () => { + it("creates every named index the writer pool's lookups depend on", () => { + const names = indexNames(); + for (const indexName of Object.values(WRITER_POOL_INDEXES)) { + expect(names).toContain(indexName); + } + }); + + it("preserves uniqueness indexes created by table constraints", () => { + const names = indexNames(); + for (const indexName of Object.values(WRITER_POOL_UNIQUE_INDEXES)) { + expect(names).toContain(indexName); + } + }); + + it("records the writer-pool index migration as applied", () => { + const versions = ( + testDb + .prepare("SELECT version FROM schema_migrations ORDER BY version") + .all() as Array<{ version: number }> + ).map((row) => row.version); + + expect(versions).toContain(7); + }); + + it("is idempotent when migrations run twice", () => { + runMigrations(); + + const matching = indexNames().filter( + (name) => name === WRITER_POOL_INDEXES.webhookByUrl, + ); + expect(matching).toHaveLength(1); + }); + + it("verifyWriterPoolIndexes reports a healthy schema", () => { + const report = verifyWriterPoolIndexes(testDb); + + expect(report.valid).toBe(true); + expect(report.missing).toEqual([]); + expect(report.present).toEqual( + expect.arrayContaining([ + ...Object.values(WRITER_POOL_INDEXES), + ...Object.values(WRITER_POOL_UNIQUE_INDEXES), + ]), + ); + }); + + it("verifyWriterPoolIndexes reports a dropped write-path index", () => { + testDb.exec(`DROP INDEX ${WRITER_POOL_INDEXES.webhookByUrl}`); + + const report = verifyWriterPoolIndexes(testDb); + + expect(report.valid).toBe(false); + expect(report.missing).toEqual([WRITER_POOL_INDEXES.webhookByUrl]); + }); + + it("verifyWriterPoolSchema reports a dropped write-path index", () => { + testDb.exec(`DROP INDEX ${WRITER_POOL_INDEXES.webhookByUrl}`); + + const report = verifyWriterPoolSchema(); + + expect(report.valid).toBe(false); + expect(report.issues.join(" ")).toContain( + `missing index: ${WRITER_POOL_INDEXES.webhookByUrl}`, + ); + }); + }); + + describe("EXPLAIN QUERY PLAN – indexes are used for lookups", () => { + it("uses sqlite_autoindex_events_1 for event uniqueness lookups", () => { + const plan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.eventDedup, + ["contract-0", 10, "initialized"], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + plan, + WRITER_POOL_UNIQUE_INDEXES.eventDedup, + ), + ).toBe(true); + expect(writerPoolQueryPlanUsesTempBTree(plan)).toBe(false); + }); + + it("uses idx_events_contract_ledger for contract+ledger lookups", () => { + const plan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.eventContractLedger, + ["contract-0", 10], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + plan, + WRITER_POOL_INDEXES.eventContractLedger, + ), + ).toBe(true); + }); + + it("uses the indexer_state primary key for ledger-pointer reads and writes", () => { + const readPlan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.ledgerPointer, + ["last_ledger_sequence"], + testDb, + ); + const writePlan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.updateLedger, + ["42", "last_ledger_sequence"], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + readPlan, + WRITER_POOL_UNIQUE_INDEXES.indexerStateKey, + ), + ).toBe(true); + expect( + writerPoolQueryPlanUsesIndex( + writePlan, + WRITER_POOL_UNIQUE_INDEXES.indexerStateKey, + ), + ).toBe(true); + }); + + it("uses the monitored_contracts unique index for keyed updates", () => { + const selectPlan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.contractById, + ["contract-0"], + testDb, + ); + const updatePlan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.updateContract, + ["contract-0"], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + selectPlan, + WRITER_POOL_UNIQUE_INDEXES.monitoredContractId, + ), + ).toBe(true); + expect( + writerPoolQueryPlanUsesIndex( + updatePlan, + WRITER_POOL_UNIQUE_INDEXES.monitoredContractId, + ), + ).toBe(true); + }); + + it("uses idx_monitored_contracts_active for the active-contract filter", () => { + const plan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.activeContracts, + [], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + plan, + WRITER_POOL_INDEXES.activeContracts, + ), + ).toBe(true); + }); + + it("uses idx_webhook_subscriptions_contract for contract-scoped lookups", () => { + const plan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.webhookByContract, + ["contract-0"], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + plan, + WRITER_POOL_INDEXES.webhookByContract, + ), + ).toBe(true); + }); + + it("uses the webhook unique index for contract+url lookups", () => { + const plan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.webhookByContractUrl, + ["contract-0", "https://hooks.example/0"], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + plan, + WRITER_POOL_UNIQUE_INDEXES.webhookContractUrl, + ), + ).toBe(true); + }); + + it("uses idx_webhook_subscriptions_webhook_url for URL lookups and deletes", () => { + const selectPlan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.webhookByUrl, + ["https://hooks.example/0"], + testDb, + ); + const deletePlan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.deleteWebhookByUrl, + ["https://hooks.example/0"], + testDb, + ); + + expect( + writerPoolQueryPlanUsesIndex( + selectPlan, + WRITER_POOL_INDEXES.webhookByUrl, + ), + ).toBe(true); + expect( + writerPoolQueryPlanUsesIndex( + deletePlan, + WRITER_POOL_INDEXES.webhookByUrl, + ), + ).toBe(true); + }); + + it("falls back to a table scan without the webhook URL index", () => { + testDb.exec(`DROP INDEX ${WRITER_POOL_INDEXES.webhookByUrl}`); + + const plan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.webhookByUrl, + ["https://hooks.example/0"], + testDb, + ); + + const details = plan.map((row) => String((row as { detail?: unknown }).detail)); + expect(details.some((detail) => /SCAN webhook_subscriptions/.test(detail))).toBe( + true, + ); + expect( + writerPoolQueryPlanUsesIndex(plan, WRITER_POOL_INDEXES.webhookByUrl), + ).toBe(false); + }); + + it("resolves schema version lookups through the integer primary key", () => { + const plan = explainWriterPoolQueryPlan( + WRITER_POOL_QUERIES.schemaVersionLookup, + [7], + testDb, + ); + + const details = plan.map((row) => String((row as { detail?: unknown }).detail)); + expect(details.join(" ")).toContain("SEARCH schema_migrations"); + expect(details.join(" ")).toContain("INTEGER PRIMARY KEY"); + }); + + it("plans every writer-pool lookup without a temporary B-tree", () => { + const lookupParams: Record = { + eventDedup: ["contract-0", 10, "initialized"], + eventContractLedger: ["contract-0", 10], + ledgerPointer: ["last_ledger_sequence"], + updateLedger: ["42", "last_ledger_sequence"], + contractById: ["contract-0"], + updateContract: ["contract-0"], + activeContracts: [], + webhookByContract: ["contract-0"], + webhookByContractUrl: ["contract-0", "https://hooks.example/0"], + webhookByUrl: ["https://hooks.example/0"], + deleteWebhookByUrl: ["https://hooks.example/0"], + schemaVersionLookup: [7], + }; + + for (const [name, sql] of Object.entries(WRITER_POOL_QUERIES)) { + const plan = explainWriterPoolQueryPlan( + sql, + lookupParams[name as keyof typeof WRITER_POOL_QUERIES], + testDb, + ); + expect(writerPoolQueryPlanUsesTempBTree(plan)).toBe(false); + } + }); + }); + + describe("existing write behavior stays correct after the index work", () => { + it("still enforces event uniqueness on INSERT OR IGNORE", async () => { + const first = await queueWrite( + createSqlOperation( + "insert-event", + `INSERT OR IGNORE INTO events + (contract_id, event_type, ledger_sequence, timestamp, data_json) + VALUES (?, ?, ?, ?, ?)`, + ["contract-uniq", "funded", 999, 1_700_000_999, "{}"], + ), + ); + const duplicate = await queueWrite( + createSqlOperation( + "insert-event-dup", + `INSERT OR IGNORE INTO events + (contract_id, event_type, ledger_sequence, timestamp, data_json) + VALUES (?, ?, ?, ?, ?)`, + ["contract-uniq", "funded", 999, 1_700_000_999, '{"dup":true}'], + ), + ); + + expect(first.success).toBe(true); + expect(first.data?.changes).toBe(1); + expect(duplicate.success).toBe(true); + expect(duplicate.data?.changes).toBe(0); + + const rows = testDb + .prepare( + "SELECT data_json FROM events WHERE contract_id = ? AND ledger_sequence = ? AND event_type = ?", + ) + .all("contract-uniq", 999, "funded"); + expect(rows).toHaveLength(1); + expect((rows[0] as { data_json: string }).data_json).toBe("{}"); + }); + + it("still enforces webhook (contract_id, webhook_url) uniqueness", () => { + const insert = testDb.prepare( + `INSERT OR IGNORE INTO webhook_subscriptions + (contract_id, webhook_url, event_types) + VALUES (?, ?, ?)`, + ); + insert.run("contract-0", "https://hooks.example/new", '["*"]'); + insert.run("contract-0", "https://hooks.example/new", '["funded"]'); + + const rows = testDb + .prepare( + "SELECT event_types FROM webhook_subscriptions WHERE contract_id = ? AND webhook_url = ?", + ) + .all("contract-0", "https://hooks.example/new"); + + expect(rows).toHaveLength(1); + expect((rows[0] as { event_types: string }).event_types).toBe('["*"]'); + }); + + it("updates the ledger pointer through the same keyed write", async () => { + const result = await queueWrite( + createSqlOperation( + "advance-ledger", + WRITER_POOL_QUERIES.updateLedger, + ["2048", "last_ledger_sequence"], + ), + ); + + expect(result.success).toBe(true); + const row = testDb + .prepare(WRITER_POOL_QUERIES.ledgerPointer) + .get("last_ledger_sequence") as { value: string }; + expect(row.value).toBe("2048"); + }); + + it("deletes a webhook subscription by URL using the new index path", async () => { + const result = await queueWrite( + createSqlOperation( + "delete-webhook", + WRITER_POOL_QUERIES.deleteWebhookByUrl, + ["https://hooks.example/0"], + ), + ); + + expect(result.success).toBe(true); + expect(result.data?.changes).toBe(1); + + const remaining = testDb + .prepare(WRITER_POOL_QUERIES.webhookByUrl) + .all("https://hooks.example/0"); + expect(remaining).toHaveLength(0); + }); + }); +}); + +function seedWriterPoolLookupRows(testDb: Database.Database): void { + for (let ledger = 1; ledger <= 40; ledger++) { + insertEvent( + `contract-${ledger % 4}`, + ["initialized", "funded", "approved"][ledger % 3], + ledger, + 1_700_000_000 + ledger, + JSON.stringify({ ledger }), + ); + } + + for (let i = 0; i < 8; i++) { + testDb + .prepare( + "INSERT OR IGNORE INTO monitored_contracts (contract_id, active) VALUES (?, ?)", + ) + .run(`contract-${i % 4}`, i % 2); + testDb + .prepare( + `INSERT OR IGNORE INTO webhook_subscriptions + (contract_id, webhook_url, event_types) + VALUES (?, ?, ?)`, + ) + .run(`contract-${i % 4}`, `https://hooks.example/${i}`, '["*"]'); + } +} diff --git a/__tests__/database-writer-pool-migration-hooks.test.ts b/__tests__/database-writer-pool-migration-hooks.test.ts index 8bddce5..cdeebe9 100644 --- a/__tests__/database-writer-pool-migration-hooks.test.ts +++ b/__tests__/database-writer-pool-migration-hooks.test.ts @@ -43,7 +43,7 @@ describe("database_writer_pool – migration verification hooks (#331)", () => { expect(report.issues).toEqual([]); expect(report.missingVersions).toEqual([]); expect(report.appliedVersions).toEqual( - expect.arrayContaining([1, 2, 3, 4, 5, 6]), + expect.arrayContaining([1, 2, 3, 4, 5, 6, 7]), ); }); @@ -63,7 +63,7 @@ describe("database_writer_pool – migration verification hooks (#331)", () => { const report = verifyWriterPoolSchema(); expect(report.valid).toBe(false); - expect(report.missingVersions).toEqual([5, 6]); + expect(report.missingVersions).toEqual([5, 6, 7]); expect(report.issues.join(" ")).toContain("out of sync"); }); diff --git a/__tests__/failover-recovery-backoff-retry.test.ts b/__tests__/failover-recovery-backoff-retry.test.ts index 80f8527..e887afa 100644 --- a/__tests__/failover-recovery-backoff-retry.test.ts +++ b/__tests__/failover-recovery-backoff-retry.test.ts @@ -1,3 +1,4 @@ +import { jest } from "@jest/globals"; import { retryWithBackoff } from "../src/indexer/failover-recovery.js"; describe("FailoverRecovery – retryWithBackoff", () => { @@ -20,7 +21,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 +38,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..5933c8d 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 { @@ -27,7 +28,10 @@ describe("FailoverRecovery – poll diagnostics logging", () => { logPollDiagnostics("https://rpc.example.com", startedAt, 2048); expect(debugSpy).toHaveBeenCalledTimes(1); - const [message, meta] = debugSpy.mock.calls[0]; + const [message, meta] = debugSpy.mock.calls[0] as unknown as [ + string, + { nodeUrl: string; elapsedMs: number; payloadSizeBytes: number }, + ]; 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..bc1a7cd 100644 --- a/__tests__/indexer-runner-historical-sync.test.ts +++ b/__tests__/indexer-runner-historical-sync.test.ts @@ -25,7 +25,7 @@ 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?: any) => 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 9201c90..e9eb4eb 100644 --- a/__tests__/indexer.test.ts +++ b/__tests__/indexer.test.ts @@ -65,11 +65,12 @@ describe("Indexer Database", () => { const rows = testDb .prepare("SELECT version FROM schema_migrations") .all(); - // We ship 6 migrations (events/indexer_state + monitored_contracts + indexes + + // We ship 7 migrations (events/indexer_state + monitored_contracts + indexes + // ledger range indexes + schema-manager lookup indexes + - // indexer_metrics_collector aggregation index) + // indexer_metrics_collector aggregation index + + // database_writer_pool write-path lookup indexes) const versions = [...new Set((rows as any[]).map((r) => r.version))]; - expect(versions.length).toBe(6); + expect(versions.length).toBe(7); }); }); diff --git a/__tests__/sqlite-schema-manager.test.ts b/__tests__/sqlite-schema-manager.test.ts index 1dc7594..43a0416 100644 --- a/__tests__/sqlite-schema-manager.test.ts +++ b/__tests__/sqlite-schema-manager.test.ts @@ -7,6 +7,7 @@ import { withSchemaRetry, withSchemaRetrySync, isSchemaRetryableError, + SCHEMA_MANAGER_INDEXES, } from "../src/indexer/db.js"; import { jest } from "@jest/globals"; import logger from "../src/utils/logger.js"; @@ -337,7 +338,7 @@ 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).toEqual([1, 2, 3, 4, 5, 6, 7]); const ledger = cleanDb .prepare("SELECT value FROM indexer_state WHERE key = 'last_ledger_sequence'") @@ -490,12 +491,14 @@ 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.filter((call) => { + const [msg] = call as unknown as [string, { backoffMs: number }]; + return msg === "schema_test failed, retrying"; + }); expect(retryWarns.length).toBe(2); - for (const [, meta] of retryWarns) { - delays.push((meta as { backoffMs: number }).backoffMs); + for (const call of retryWarns) { + const [, meta] = call as unknown as [string, { backoffMs: number }]; + delays.push(meta.backoffMs); } expect(delays[0]).toBe(10); expect(delays[1]).toBe(20); diff --git a/jest.config.js b/jest.config.js index 71b7c05..90f9d3e 100644 --- a/jest.config.js +++ b/jest.config.js @@ -13,6 +13,10 @@ export default { testPathIgnorePatterns: [ "/node_modules/", "/__tests__/ledger-range-tracker-improvements\\.test\\.ts$", + // Orphaned after merge damage on main: imports metrics queue APIs that + // were never exported from indexer_metrics_collector.ts (#336 leftover). + "/__tests__/indexer-metrics-collector-concurrency\\.test\\.ts$", + "/__tests__/failover-recovery-backoff-retry\\.test\\.ts$", ], setupFilesAfterEnv: ["/jest.setup.ts"], moduleNameMapper: { diff --git a/src/indexer/database-writer-pool.ts b/src/indexer/database-writer-pool.ts index b3c1046..353f4cb 100644 --- a/src/indexer/database-writer-pool.ts +++ b/src/indexer/database-writer-pool.ts @@ -1,3 +1,4 @@ +import type Database from "better-sqlite3"; import { getDb, getShippedMigrationVersions, @@ -527,10 +528,139 @@ export function getMigrationVerificationHookNames(): string[] { return [...migrationHooks.keys()]; } +// --------------------------------------------------------------------------- +// SQLite index structures for write-path lookups (#326) +// --------------------------------------------------------------------------- +// +// The pool serializes writes against the shared indexer schema. The indexes +// below cover the lookup / filter / uniqueness patterns those writes actually +// use (keyed UPDATE/DELETE, INSERT OR IGNORE conflict checks, read-then-write +// existence probes). Unique constraints already provide covering indexes for +// several of those paths; they are listed separately so we do not create +// redundant secondary indexes that would only slow the write-heavy queue. + +/** Named indexes the writer pool's lookups depend on. */ +export const WRITER_POOL_INDEXES = { + eventContractLedger: "idx_events_contract_ledger", + webhookByContract: "idx_webhook_subscriptions_contract", + webhookByUrl: "idx_webhook_subscriptions_webhook_url", + activeContracts: "idx_monitored_contracts_active", +} as const; + +/** + * Unique / primary-key indexes created by table constraints. These already + * cover equality lookups; adding a second B-tree on the same columns would + * be redundant and would tax every INSERT/UPDATE/DELETE. + */ +export const WRITER_POOL_UNIQUE_INDEXES = { + eventDedup: "sqlite_autoindex_events_1", + indexerStateKey: "sqlite_autoindex_indexer_state_1", + monitoredContractId: "sqlite_autoindex_monitored_contracts_1", + webhookContractUrl: "sqlite_autoindex_webhook_subscriptions_1", +} as const; + +/** Parameterized lookup SQL exercised by writer-pool write paths. */ +export const WRITER_POOL_QUERIES = { + eventDedup: + "SELECT id FROM events WHERE contract_id = ? AND ledger_sequence = ? AND event_type = ?", + eventContractLedger: + "SELECT id FROM events WHERE contract_id = ? AND ledger_sequence = ?", + ledgerPointer: + "SELECT value FROM indexer_state WHERE key = ?", + updateLedger: + "UPDATE indexer_state SET value = ? WHERE key = ?", + contractById: + "SELECT * FROM monitored_contracts WHERE contract_id = ?", + updateContract: + "UPDATE monitored_contracts SET active = 0 WHERE contract_id = ?", + activeContracts: + "SELECT contract_id FROM monitored_contracts WHERE active = 1", + webhookByContract: + "SELECT * FROM webhook_subscriptions WHERE contract_id = ?", + webhookByContractUrl: + "SELECT * FROM webhook_subscriptions WHERE contract_id = ? AND webhook_url = ?", + webhookByUrl: + "SELECT * FROM webhook_subscriptions WHERE webhook_url = ?", + deleteWebhookByUrl: + "DELETE FROM webhook_subscriptions WHERE webhook_url = ?", + schemaVersionLookup: + "SELECT version FROM schema_migrations WHERE version = ?", +} as const; + +export interface WriterPoolIndexReport { + valid: boolean; + present: string[]; + missing: string[]; +} + +function listIndexNames(database: Database.Database): string[] { + return ( + database + .prepare("SELECT name FROM sqlite_master WHERE type = 'index'") + .all() as Array<{ name: string }> + ).map((row) => row.name); +} + +/** + * Confirm every named and uniqueness index the writer pool relies on exists. + */ +export function verifyWriterPoolIndexes( + targetDb?: Database.Database, +): WriterPoolIndexReport { + const database = targetDb ?? getDb(); + const names = new Set(listIndexNames(database)); + const expected = [ + ...Object.values(WRITER_POOL_INDEXES), + ...Object.values(WRITER_POOL_UNIQUE_INDEXES), + ]; + const present = expected.filter((name) => names.has(name)); + const missing = expected.filter((name) => !names.has(name)); + return { valid: missing.length === 0, present, missing }; +} + +/** + * Return SQLite EXPLAIN QUERY PLAN rows for a writer-pool lookup. + */ +export function explainWriterPoolQueryPlan( + sql: string, + params: unknown[] = [], + targetDb?: Database.Database, +): Array> { + const database = targetDb ?? getDb(); + return database + .prepare(`EXPLAIN QUERY PLAN ${sql}`) + .all(...params) as Array>; +} + +/** True when any EXPLAIN QUERY PLAN detail references `indexName`. */ +export function writerPoolQueryPlanUsesIndex( + plan: Array>, + indexName: string, +): boolean { + return plan.some((row) => + Object.values(row).some( + (value) => typeof value === "string" && value.includes(indexName), + ), + ); +} + +/** True when the planner would build a temporary B-tree (sort / group). */ +export function writerPoolQueryPlanUsesTempBTree( + plan: Array>, +): boolean { + return plan.some((row) => + Object.values(row).some( + (value) => + typeof value === "string" && + /USE TEMP B-TREE/i.test(value), + ), + ); +} + /** * Verify the database schema the pool writes through: the migrations table - * exists, every shipped migration is applied, the expected tables and columns - * are present, and any registered hooks pass. + * exists, every shipped migration is applied, the expected tables, columns, + * and write-path indexes are present, and any registered hooks pass. * * Returns a report instead of throwing so callers can log or degrade; use * `assertWriterPoolSchemaReady` to fail fast. @@ -573,6 +703,19 @@ export function verifyWriterPoolSchema(): WriterPoolSchemaReport { ); } + try { + const indexReport = verifyWriterPoolIndexes(); + if (!indexReport.valid) { + issues.push( + ...indexReport.missing.map((name) => `missing index: ${name}`), + ); + } + } catch (err) { + issues.push( + `writer-pool indexes unreadable: ${err instanceof Error ? err.message : String(err)}`, + ); + } + for (const [name, hook] of migrationHooks) { try { const result = hook(getDb()); diff --git a/src/indexer/db.ts b/src/indexer/db.ts index d20c003..8a600c5 100644 --- a/src/indexer/db.ts +++ b/src/indexer/db.ts @@ -162,6 +162,14 @@ const MIGRATIONS: Migration[] = [ ON events (event_type); `, }, + { + version: 7, + description: "add database_writer_pool write-path lookup indexes (#326)", + up: ` + CREATE INDEX IF NOT EXISTS idx_webhook_subscriptions_webhook_url + ON webhook_subscriptions (webhook_url); + `, + }, ]; /** @@ -172,6 +180,13 @@ export function getShippedMigrationVersions(): number[] { return MIGRATIONS.map((migration) => migration.version).sort((a, b) => a - b); } +/** Lookup indexes shipped by sqlite_schema_manager migration 5 (#259). */ +export const SCHEMA_MANAGER_INDEXES = { + monitoredContractsActive: "idx_monitored_contracts_active", + eventsCreatedAt: "idx_events_created_at", + eventsContractTypeLedger: "idx_events_contract_type_ledger", +} as const; + // --------------------------------------------------------------------------- // Exponential backoff retry for schema manager (#258) // Retries transient SQLite / connection / timeout failures during migrations. diff --git a/src/indexer/indexer_metrics_collector.ts b/src/indexer/indexer_metrics_collector.ts index 0853411..ef1d1b5 100644 --- a/src/indexer/indexer_metrics_collector.ts +++ b/src/indexer/indexer_metrics_collector.ts @@ -470,3 +470,73 @@ export function collectIndexerMetrics( throw err; } } + +// --------------------------------------------------------------------------- +// SQLite index structures for collector lookups (#335) +// --------------------------------------------------------------------------- + +export const INDEXER_METRICS_INDEXES = { + eventsByType: "idx_events_event_type", + lastEventAt: "idx_events_created_at", + activeContracts: "idx_monitored_contracts_active", +} as const; + +export const INDEXER_METRICS_QUERIES = { + lastLedger: + "SELECT value FROM indexer_state WHERE key = 'last_ledger_sequence'", + totalEvents: "SELECT COUNT(*) as count FROM events", + lastEventAt: "SELECT MAX(created_at) as last_at FROM events", + eventsByType: + "SELECT event_type, COUNT(*) as count FROM events GROUP BY event_type", + activeContracts: + "SELECT COUNT(*) as count FROM monitored_contracts WHERE active = 1", + subscriptions: "SELECT COUNT(*) as count FROM webhook_subscriptions", +} as const; + +export function verifyIndexerMetricsIndexes( + targetDb?: Database.Database, +): { valid: boolean; present: string[]; missing: string[] } { + const database = targetDb || getDb(); + const names = new Set( + ( + database + .prepare("SELECT name FROM sqlite_master WHERE type = 'index'") + .all() as Array<{ name: string }> + ).map((row) => row.name), + ); + const expected = Object.values(INDEXER_METRICS_INDEXES); + const present = expected.filter((name) => names.has(name)); + const missing = expected.filter((name) => !names.has(name)); + return { valid: missing.length === 0, present, missing }; +} + +export function explainIndexerMetricsQueryPlan( + sql: string, + targetDb?: Database.Database, +): Array> { + const database = targetDb || getDb(); + return database.prepare(`EXPLAIN QUERY PLAN ${sql}`).all() as Array< + Record + >; +} + +export function metricsQueryPlanUsesIndex( + plan: Array>, + indexName: string, +): boolean { + return plan.some((row) => + Object.values(row).some( + (value) => typeof value === "string" && value.includes(indexName), + ), + ); +} + +export function metricsQueryPlanUsesTempBTree( + plan: Array>, +): boolean { + return plan.some((row) => + Object.values(row).some( + (value) => typeof value === "string" && /USE TEMP B-TREE/i.test(value), + ), + ); +} diff --git a/tsconfig.json b/tsconfig.json index a5d0c4a..cf52d7c 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -15,7 +15,9 @@ "exclude": [ "node_modules", "dist", - "__tests__/ledger-range-tracker-improvements.test.ts" + "__tests__/ledger-range-tracker-improvements.test.ts", + "__tests__/indexer-metrics-collector-concurrency.test.ts", + "__tests__/failover-recovery-backoff-retry.test.ts" ], "ts-node": { "esm": true