From cf694a96aa24a894927de6265f5c47572c886df0 Mon Sep 17 00:00:00 2001 From: Adam Spitz Date: Mon, 28 Sep 2026 11:49:22 -0400 Subject: [PATCH] Replay the testnet chain head while the indexer is idle. Base Sepolia keeps Ponder's four-second timer but serves the last latest-block response for 30s unless a client read or POST /api/indexer-wake opens a one-minute live window. --- indexer/README.md | 2 +- indexer/package.json | 2 +- indexer/src/api/index.ts | 14 ++ indexer/src/indexing/ponderEnv.ts | 2 + indexer/src/rpc/idleHeadCache.test.ts | 137 ++++++++++++++++ indexer/src/rpc/idleHeadCache.ts | 227 ++++++++++++++++++++++++++ workflow/deployment.md | 2 +- 7 files changed, 383 insertions(+), 3 deletions(-) create mode 100644 indexer/src/rpc/idleHeadCache.test.ts create mode 100644 indexer/src/rpc/idleHeadCache.ts diff --git a/indexer/README.md b/indexer/README.md index 48f13354..3ee9cad1 100644 --- a/indexer/README.md +++ b/indexer/README.md @@ -48,7 +48,7 @@ For Render or other hosted environments: start mode, and Render's native file-watcher limit can otherwise abort startup with `EMFILE: too many open files, watch '/app'`. - Keep `PONDER_ETH_GET_LOGS_BLOCK_RANGE` large enough for catch-up. The Render blueprint defaults to `10000`. A tiny range such as `10` makes a million-block historical sync require hundreds of thousands of `eth_getLogs` batches and will blow Alchemy CUPS. If the provider rejects the window, the process logs a one-shot `[commonality-indexer] eth_getLogs failed because the RPC rejected the block range or response size` line with the env to change; lower to `1000` then `10`. -- Hosted chains poll every `PONDER_POLL_INTERVAL_MS` (default `4000`). Hardhat stays at 100ms. +- Hosted chains poll every `PONDER_POLL_INTERVAL_MS` (default `4000`). Hardhat stays at 100ms. On Base Sepolia the process still ticks that often, but while nobody is reading it replays the last `eth_getBlockByNumber("latest")` for `INDEXER_IDLE_HEAD_INTERVAL_MS` (default 30s) instead of calling Alchemy. A GraphQL or `/api` read, or `POST /api/indexer-wake` after a submitted transaction, passes polls through for `INDEXER_IDLE_HEAD_WAKE_MS` (default 60s). Health checks and `/api/project-read-demand` do not wake it. Set `INDEXER_IDLE_HEAD_CACHE=0` to disable, or `=1` to enable on mainnet. - Keep `DATABASE_SCHEMA` stable (`commonality_base_sepolia_v6` on testnet). Renaming it drops the event cache and replays history against the RPC. `scripts/smoke-check-render.mjs` fails if the name changes unless `INDEXER_ALLOW_SCHEMA_BUMP=1`. Code deploys (`ponder start` + `PONDER_EXPERIMENTAL_DB=platform` + the persistent disk for stop-before-start) reuse the same schema. The `events` table is append-only raw logs; new handlers and extra contract addresses in `INDEXER_DEPLOYMENT_MANIFEST` do not need a wipe. ### RPC budget (Alchemy monthly CU) diff --git a/indexer/package.json b/indexer/package.json index a0540707..233e026d 100644 --- a/indexer/package.json +++ b/indexer/package.json @@ -14,7 +14,7 @@ "serve": "ponder serve", "lint": "eslint .", "typecheck": "tsc --noEmit && npm run check-abis", - "test": "node --import tsx --test selectPonderScript.test.mjs src/rpc/ethGetLogsRangeGuard.test.ts src/rpc/monthlyCapacity.test.ts src/api/projectReadDemand.test.ts src/indexing/contractCapabilities.test.ts src/indexing/conceptspaceConfigGraph.test.ts", + "test": "node --import tsx --test selectPonderScript.test.mjs src/rpc/ethGetLogsRangeGuard.test.ts src/rpc/monthlyCapacity.test.ts src/rpc/idleHeadCache.test.ts src/api/projectReadDemand.test.ts src/indexing/contractCapabilities.test.ts src/indexing/conceptspaceConfigGraph.test.ts", "build": "tsc", "check-abis": "tsx scripts/sync-abis.ts --check", "clean": "rm -rf .ponder", diff --git a/indexer/src/api/index.ts b/indexer/src/api/index.ts index 4d75d2c9..a252e929 100644 --- a/indexer/src/api/index.ts +++ b/indexer/src/api/index.ts @@ -14,6 +14,7 @@ import { client, graphql } from "ponder"; import { and, desc, eq, gte, lte, or } from "ponder"; import { getAddress, isAddress, type Hex } from "viem"; import { fundingIndexerRoutesEnabled } from "../indexing/contractCapabilities"; +import { noteIndexerWake, requestWakesIndexer } from "../rpc/idleHeadCache"; import { isBareContractLogQuery, projectReadDemandReport, recordUnindexedProjectLogRequest } from "./projectReadDemand"; /** @@ -80,6 +81,19 @@ function publicationPointer(event: { blockNumber: bigint; transactionHash: strin const app = new Hono(); +app.use("*", async (c, next) => { + if (requestWakesIndexer(c.req.method, c.req.path)) noteIndexerWake(); + await next(); +}); + +app.post("/api/indexer-wake", (c) => { + const until = noteIndexerWake(); + return c.json({ + ok: true, + fastPollingUntil: until > 0 ? new Date(until).toISOString() : null, + }); +}); + // Expose SQL client for direct queries (all tables) app.use("/sql/*", client({ db, schema })); diff --git a/indexer/src/indexing/ponderEnv.ts b/indexer/src/indexing/ponderEnv.ts index da4f0708..922d153f 100644 --- a/indexer/src/indexing/ponderEnv.ts +++ b/indexer/src/indexing/ponderEnv.ts @@ -1,5 +1,6 @@ import { http } from "viem"; import { installEthGetLogsRangeGuard } from "../rpc/ethGetLogsRangeGuard"; +import { idleHeadCacheEnabled, installIdleHeadCache } from "../rpc/idleHeadCache"; import { installMonthlyCapacityGuard } from "../rpc/monthlyCapacity"; import { INDEXER_CHAIN_IDS, type IndexerChainName } from "../utils/chain"; @@ -222,4 +223,5 @@ export function installHostedRpcGuards(context: IndexerDeploymentContext): void : context.ethGetLogsBlockRange; installEthGetLogsRangeGuard({ configuredRange }); installMonthlyCapacityGuard(); + if (idleHeadCacheEnabled(context.chain)) installIdleHeadCache(); } diff --git a/indexer/src/rpc/idleHeadCache.test.ts b/indexer/src/rpc/idleHeadCache.test.ts new file mode 100644 index 00000000..ea41e218 --- /dev/null +++ b/indexer/src/rpc/idleHeadCache.test.ts @@ -0,0 +1,137 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { + createIdleHeadCache, + idleHeadCacheEnabled, + requestWakesIndexer, +} from "./idleHeadCache"; + +const head = (id: string, number = "0x10") => + JSON.stringify({ jsonrpc: "2.0", id, result: { number, hash: "0xabc" } }); + +test("idle cache is on for base sepolia unless disabled", () => { + assert.equal(idleHeadCacheEnabled("base-sepolia", {}), true); + assert.equal(idleHeadCacheEnabled("base-sepolia", { INDEXER_IDLE_HEAD_CACHE: "0" }), false); + assert.equal(idleHeadCacheEnabled("mainnet", {}), false); + assert.equal(idleHeadCacheEnabled("mainnet", { INDEXER_IDLE_HEAD_CACHE: "1" }), true); +}); + +test("client reads and the wake endpoint count; health and demand do not", () => { + assert.equal(requestWakesIndexer("POST", "/graphql"), true); + assert.equal(requestWakesIndexer("GET", "/api/events"), true); + assert.equal(requestWakesIndexer("POST", "/api/indexer-wake"), true); + assert.equal(requestWakesIndexer("GET", "/api/indexer-wake"), false); + assert.equal(requestWakesIndexer("GET", "/health"), false); + assert.equal(requestWakesIndexer("GET", "/ready"), false); + assert.equal(requestWakesIndexer("GET", "/api/project-read-demand"), false); +}); + +test("replays the latest block while idle and refreshes after the interval", async () => { + let clock = 1_000; + const calls: string[] = []; + const fetchImpl: typeof fetch = async (_input, init) => { + const body = String(init?.body); + calls.push(body); + const id = JSON.parse(body).id; + return new Response(head(id, calls.length === 1 ? "0x10" : "0x11"), { status: 200 }); + }; + const cache = createIdleHeadCache({ idleIntervalMs: 30_000, wakeMs: 60_000, now: () => clock, fetchImpl }); + const restore = cache.install(); + try { + const first = await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "eth_getBlockByNumber", params: ["latest", true] }), + }); + assert.equal(JSON.parse(await first.text()).result.number, "0x10"); + + clock = 10_000; + const second = await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 2, method: "eth_getBlockByNumber", params: ["latest", true] }), + }); + const replayed = JSON.parse(await second.text()); + assert.equal(replayed.id, 2); + assert.equal(replayed.result.number, "0x10"); + assert.equal(calls.length, 1); + + clock = 31_000; + const third = await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 3, method: "eth_getBlockByNumber", params: ["latest", true] }), + }); + assert.equal(JSON.parse(await third.text()).result.number, "0x11"); + assert.equal(calls.length, 2); + } finally { + restore(); + } +}); + +test("a wake lets the next poll through, then idle caching resumes", async () => { + let clock = 5_000; + let upstream = 0; + const fetchImpl: typeof fetch = async () => { + upstream += 1; + return new Response(head(String(upstream), `0x${upstream.toString(16)}`), { status: 200 }); + }; + const cache = createIdleHeadCache({ idleIntervalMs: 30_000, wakeMs: 60_000, now: () => clock, fetchImpl }); + const restore = cache.install(); + const body = JSON.stringify({ jsonrpc: "2.0", id: 1, method: "eth_getBlockByNumber", params: ["latest", true] }); + try { + await fetch("https://rpc.example", { method: "POST", body }); + clock = 6_000; + await fetch("https://rpc.example", { method: "POST", body }); + assert.equal(upstream, 1); + + cache.noteWake(); + clock = 7_000; + await fetch("https://rpc.example", { method: "POST", body }); + assert.equal(upstream, 2); + + clock = 66_000; + await fetch("https://rpc.example", { method: "POST", body }); + assert.equal(upstream, 3); + + clock = 80_000; + await fetch("https://rpc.example", { method: "POST", body }); + assert.equal(upstream, 3); + } finally { + restore(); + } +}); + +test("does not cache logs, historical blocks, or errors", async () => { + let upstream = 0; + const fetchImpl: typeof fetch = async (_input, init) => { + upstream += 1; + const method = JSON.parse(String(init?.body)).method; + if (method === "eth_getLogs") return new Response(JSON.stringify({ jsonrpc: "2.0", id: 1, result: [] }), { status: 200 }); + return new Response(JSON.stringify({ jsonrpc: "2.0", id: 1, error: { code: -32000, message: "nope" } }), { status: 200 }); + }; + const cache = createIdleHeadCache({ now: () => 1_000, fetchImpl }); + const restore = cache.install(); + try { + await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "eth_getLogs", params: [{}] }), + }); + await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "eth_getLogs", params: [{}] }), + }); + await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "eth_getBlockByNumber", params: ["0x10", false] }), + }); + await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "eth_getBlockByNumber", params: ["latest", true] }), + }); + await fetch("https://rpc.example", { + method: "POST", + body: JSON.stringify({ jsonrpc: "2.0", id: 1, method: "eth_getBlockByNumber", params: ["latest", true] }), + }); + assert.equal(upstream, 5); + } finally { + restore(); + } +}); diff --git a/indexer/src/rpc/idleHeadCache.ts b/indexer/src/rpc/idleHeadCache.ts new file mode 100644 index 00000000..bd28b8ac --- /dev/null +++ b/indexer/src/rpc/idleHeadCache.ts @@ -0,0 +1,227 @@ +/** + * Ponder's poll timer is fixed at startup. While idle, replay the last + * successful `eth_getBlockByNumber("latest")` so that timer does not call + * Alchemy. A reader or an explicit wake lets polls through for a short window; + * the next real head makes Ponder catch the gap with one eth_getLogs. + */ + +export const DEFAULT_IDLE_HEAD_INTERVAL_MS = 30_000; +export const DEFAULT_IDLE_HEAD_WAKE_MS = 60_000; + +type JsonRpc = { + id?: unknown; + method?: unknown; + params?: unknown; + error?: unknown; + result?: unknown; +}; + +export type IdleHeadCacheOptions = { + enabled?: boolean; + idleIntervalMs?: number; + wakeMs?: number; + now?: () => number; + fetchImpl?: typeof fetch; +}; + +export function idleHeadCacheEnabled(chain: string, env: NodeJS.ProcessEnv = process.env): boolean { + const flag = env.INDEXER_IDLE_HEAD_CACHE; + if (flag === "0" || flag === "false") return false; + if (flag === "1" || flag === "true") return true; + return chain === "base-sepolia"; +} + +export function requestWakesIndexer(method: string, path: string): boolean { + const normalized = path.split("?")[0] ?? path; + if (method !== "GET" && method !== "POST") return false; + if (normalized === "/health" || normalized === "/ready" || normalized === "/metrics") return false; + if (normalized === "/api/project-read-demand") return false; + if (normalized === "/api/indexer-wake") return method === "POST"; + if (normalized === "/" || normalized === "/graphql" || normalized.startsWith("/graphql/")) return true; + if (normalized.startsWith("/api/")) return method === "GET"; + if (normalized.startsWith("/sql")) return true; + return false; +} + +function readPositiveInt(raw: string | undefined, fallback: number): number { + if (raw === undefined || raw === "") return fallback; + const parsed = Number(raw); + if (!Number.isFinite(parsed) || parsed < 500) { + throw new Error(`Invalid idle-head duration "${raw}". Expected a millisecond count >= 500.`); + } + return parsed; +} + +function latestBlockRequests(body: string): JsonRpc[] | null { + let parsed: unknown; + try { + parsed = JSON.parse(body); + } catch { + return null; + } + const calls = Array.isArray(parsed) ? parsed : [parsed]; + if (calls.length === 0) return null; + const requests: JsonRpc[] = []; + for (const call of calls) { + if (!call || typeof call !== "object") return null; + const rpc = call as JsonRpc; + if (rpc.method !== "eth_getBlockByNumber") return null; + if (!Array.isArray(rpc.params) || rpc.params[0] !== "latest") return null; + requests.push(rpc); + } + return requests; +} + +function rewriteIds(cachedBody: string, requests: JsonRpc[]): string | null { + let parsed: unknown; + try { + parsed = JSON.parse(cachedBody); + } catch { + return null; + } + const responses = Array.isArray(parsed) ? parsed : [parsed]; + if (responses.length !== requests.length) return null; + const rewritten = responses.map((response, index) => { + if (!response || typeof response !== "object") return response; + return { ...(response as JsonRpc), id: requests[index]?.id }; + }); + return JSON.stringify(Array.isArray(parsed) ? rewritten : rewritten[0]); +} + +function cacheableHeadResponse(body: string): boolean { + let parsed: unknown; + try { + parsed = JSON.parse(body); + } catch { + return false; + } + const responses = Array.isArray(parsed) ? parsed : [parsed]; + return responses.every((response) => { + if (!response || typeof response !== "object") return false; + const rpc = response as JsonRpc; + return rpc.error === undefined && rpc.result !== undefined && rpc.result !== null; + }); +} + +export type IdleHeadCache = { + noteWake: (at?: number) => number; + fastPollingUntil: () => number; + reset: () => void; +}; + +export function createIdleHeadCache(options: IdleHeadCacheOptions = {}): IdleHeadCache & { + install: () => () => void; +} { + const idleIntervalMs = options.idleIntervalMs ?? DEFAULT_IDLE_HEAD_INTERVAL_MS; + const wakeMs = options.wakeMs ?? DEFAULT_IDLE_HEAD_WAKE_MS; + const now = options.now ?? Date.now; + let wakeUntil = 0; + let cachedBody = ""; + let cachedKey = ""; + let cachedAt = 0; + + const cache: IdleHeadCache = { + noteWake(at = now()) { + wakeUntil = at + wakeMs; + return wakeUntil; + }, + fastPollingUntil() { + return wakeUntil; + }, + reset() { + wakeUntil = 0; + cachedBody = ""; + cachedKey = ""; + cachedAt = 0; + }, + }; + + function install(): () => void { + const previous = globalThis.fetch; + const currentFetch: typeof fetch = options.fetchImpl + ? (input, init) => options.fetchImpl!(input, init) + : previous.bind(globalThis); + + const wrapped: typeof fetch = async (input, init) => { + const requestText = await peekRequestBody(input, init); + const requests = requestText ? latestBlockRequests(requestText) : null; + const key = requests ? JSON.stringify(requests.map((request) => request.params)) : ""; + const at = now(); + if ( + requests && + cachedBody && + key === cachedKey && + at >= cachedAt && + at - cachedAt < idleIntervalMs && + at >= wakeUntil + ) { + const body = rewriteIds(cachedBody, requests); + if (body) { + return new Response(body, { + status: 200, + headers: { "content-type": "application/json" }, + }); + } + } + + const response = await currentFetch(input, init); + if (!requests || !response.ok) return response; + try { + const responseText = await response.clone().text(); + if (!cacheableHeadResponse(responseText)) return response; + cachedBody = responseText; + cachedKey = key; + cachedAt = at; + } catch { + // Never break RPC on a cache write. + } + return response; + }; + + globalThis.fetch = wrapped; + return () => { + if (globalThis.fetch === wrapped) globalThis.fetch = previous; + }; + } + + return { ...cache, install }; +} + +let active: ReturnType | null = null; + +export function noteIndexerWake(at?: number): number { + if (!active) return 0; + return active.noteWake(at); +} + +export function indexerFastPollingUntil(): number { + return active?.fastPollingUntil() ?? 0; +} + +export function installIdleHeadCache(options: IdleHeadCacheOptions = {}): () => void { + if (options.enabled === false) return () => {}; + const idleIntervalMs = options.idleIntervalMs ?? readPositiveInt(process.env.INDEXER_IDLE_HEAD_INTERVAL_MS, DEFAULT_IDLE_HEAD_INTERVAL_MS); + const wakeMs = options.wakeMs ?? readPositiveInt(process.env.INDEXER_IDLE_HEAD_WAKE_MS, DEFAULT_IDLE_HEAD_WAKE_MS); + active = createIdleHeadCache({ ...options, idleIntervalMs, wakeMs }); + console.error( + `[commonality-indexer] idle head cache on: replay latest block for ${idleIntervalMs}ms, fast poll for ${wakeMs}ms after a client read or POST /api/indexer-wake.`, + ); + const restoreFetch = active.install(); + return () => { + restoreFetch(); + active = null; + }; +} + +async function peekRequestBody(input: Parameters[0], init?: RequestInit): Promise { + if (typeof init?.body === "string") return init.body; + if (init?.body instanceof Uint8Array) return new TextDecoder().decode(init.body); + if (input instanceof Request) { + try { + return await input.clone().text(); + } catch { + return ""; + } + } + return ""; +} diff --git a/workflow/deployment.md b/workflow/deployment.md index 1aa20ef2..a3c9b86c 100644 --- a/workflow/deployment.md +++ b/workflow/deployment.md @@ -433,7 +433,7 @@ The blueprint already wires: - `DATABASE_SCHEMA=commonality_base_sepolia_v6` for the current Base Sepolia deployment. **Keep this name.** Changing it wipes the event cache and forces a full historical `eth_getLogs` replay (this is how we accidentally spent a month of Alchemy CUs). Bump only when `indexer/schemas/events.schema.ts` is incompatible, and only with `INDEXER_ALLOW_SCHEMA_BUMP=1` so `npm run smoke-check` will pass. Schema-lock errors are a rolling-deploy problem — keep the persistent disk — not a reason to mint `v7`. - `PONDER_EXPERIMENTAL_DB=platform` so normal Render redeploys of a changed Ponder build can reuse the same production schema instead of failing with "previously used by a different Ponder app". - `PONDER_ETH_GET_LOGS_BLOCK_RANGE=10000` for Base Sepolia on the current Alchemy PAYG key (wider windows cut catch-up `eth_getLogs` count by ~1000× vs a 10-block free-tier cap). If logs show `[commonality-indexer] eth_getLogs failed because the RPC rejected the block range or response size`, lower it to `1000` then `10`, PUT the Render env, and **deploy** (not restart only). The indexer wraps `fetch` to print that hint once; do not switch the RPC to a viem `http()` transport to “see errors” — that breaks Ponder’s rate limiter. -- `PONDER_POLL_INTERVAL_MS=4000` on hosted chains so head-following does not poll faster than Base Sepolia block time. +- `PONDER_POLL_INTERVAL_MS=4000` on hosted chains so head-following does not poll faster than Base Sepolia block time. Base Sepolia additionally replays the cached chain head for 30s while idle (`INDEXER_IDLE_HEAD_CACHE`, on unless set to `0`). Client reads and `POST /api/indexer-wake` open a 60s fast-poll window. Do not point a health check at `/graphql` or that window stays open. - Monthly Alchemy capacity 429s: do not crash-loop and do not switch to `sepolia.base.org`. `indexer/start.sh` backs off (1m…6h) with a stub `/graphql`. Raise the dashboard usage limit, then the next retry resumes the same `DATABASE_SCHEMA`. - `START_BLOCK` is the fallback earliest block for contracts that do not set their own `startBlock`. Newly deployed contracts should use **their deploy block**, not a replay from the original 42768673. Do not lower global `START_BLOCK` or mint a new `DATABASE_SCHEMA` to “include more history.” - The indexer declares a small persistent disk even though it does not store application data there. This is an intentional Render workaround, not indexer storage: Render disables zero-downtime/rolling deploys for services with disks, which gives Ponder the stop-before-start deployment behavior it needs for the exclusive `DATABASE_SCHEMA` lock. Do not remove this disk just because `/data` appears unused unless the indexer has moved to a cleaner singleton-writer deployment model.