Skip to content

Commit 295a604

Browse files
authored
Merge pull request #756 from AugurProject/t3code/recover-pruned-log-history
Recover Augurscan from pruned RPC log history
2 parents f58e09b + 36e0d4e commit 295a604

7 files changed

Lines changed: 852 additions & 89 deletions

File tree

augurScan/src/database.ts

Lines changed: 113 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -284,6 +284,59 @@ export const lockLiveEventWriter = async (sql: SQL): Promise<void> => {
284284
await sql`SELECT singleton FROM live_event_state WHERE singleton FOR UPDATE`
285285
}
286286

287+
const canonicalHistoryTables = [
288+
'blocks',
289+
'transactions',
290+
'logs',
291+
'contract_discoveries',
292+
'questions',
293+
'pools',
294+
'pool_snapshots',
295+
'pool_state_events',
296+
'vault_snapshots',
297+
'universe_events',
298+
'amm_markets',
299+
'amm_price_snapshots',
300+
'rep_eth_price_snapshots',
301+
'uniswap_rep_eth_markets',
302+
'uniswap_rep_eth_price_observations',
303+
'protocol_timeline_entries',
304+
'open_oracle_report_events',
305+
'escalation_game_events',
306+
'truth_auction_events',
307+
'amm_trade_events',
308+
'fork_migration_events',
309+
'liquidation_approval_events',
310+
'address_activity',
311+
'address_balance_snapshots',
312+
'token_metadata',
313+
] as const
314+
315+
const invalidateCanonicalHistory = async (transaction: TransactionSQL, chainId: number, discoveryRetirementFloor?: bigint): Promise<void> => {
316+
for (const table of canonicalHistoryTables) await transaction.unsafe(`UPDATE ${table} SET canonical = false WHERE chain_id = $1 AND canonical`, [chainId])
317+
await transaction`UPDATE entity_state_snapshots SET read_status = 'stale', canonical = false WHERE chain_id = ${chainId} AND canonical`
318+
await transaction`UPDATE blocks SET finalized = false WHERE chain_id = ${chainId}`
319+
await transaction`UPDATE logs SET finalized = false WHERE chain_id = ${chainId}`
320+
await transaction`DELETE FROM log_scan_cursors WHERE chain_id = ${chainId}`
321+
await transaction`
322+
UPDATE contracts SET deployment_block = NULL, deployment_timestamp = NULL,
323+
deployment_block_exact = NULL, deployment_checked_block = NULL
324+
WHERE chain_id = ${chainId}
325+
`
326+
if (discoveryRetirementFloor === undefined)
327+
await transaction`
328+
UPDATE contracts SET canonical = false
329+
WHERE chain_id = ${chainId} AND provenance <> 'manifest'
330+
AND (discovery_block IS NULL OR discovery_block >= (SELECT start_block FROM networks WHERE chain_id = ${chainId}))
331+
`
332+
else
333+
await transaction`
334+
UPDATE contracts SET canonical = false
335+
WHERE chain_id = ${chainId} AND provenance <> 'manifest'
336+
AND (discovery_block IS NULL OR discovery_block >= ${discoveryRetirementFloor.toString()})
337+
`
338+
}
339+
287340
export const releaseReservedConnection = async (connection: Pick<ReservedSQL, 'release'>): Promise<void> => {
288341
await connection.release()
289342
}
@@ -571,49 +624,36 @@ export class ScannerDatabase {
571624
AND contract.address = discovery.address
572625
AND NOT contract.canonical
573626
`
627+
await transaction`
628+
UPDATE contracts AS contract SET
629+
label = discovery.label,
630+
kind = discovery.kind,
631+
provenance = discovery.provenance,
632+
canonical = true
633+
FROM (
634+
SELECT DISTINCT ON (candidate.address)
635+
candidate.address, candidate.label, candidate.kind, candidate.provenance
636+
FROM contract_discoveries AS candidate
637+
JOIN contracts AS retained
638+
ON retained.chain_id = candidate.chain_id
639+
AND retained.address = candidate.address
640+
AND retained.discovery_block = candidate.block_number
641+
AND retained.discovery_tx_hash = candidate.tx_hash
642+
JOIN networks AS network ON network.chain_id = candidate.chain_id
643+
WHERE candidate.chain_id = ${network.chainId} AND candidate.block_number < network.start_block
644+
ORDER BY candidate.address, candidate.canonical DESC, candidate.block_hash
645+
) AS discovery
646+
WHERE contract.chain_id = ${network.chainId}
647+
AND contract.address = discovery.address
648+
AND contract.provenance = 'manifest'
649+
AND NOT contract.canonical
650+
`
574651
await transaction`
575652
UPDATE contracts SET provenance = 'retired-manifest'
576653
WHERE chain_id = ${network.chainId} AND provenance = 'manifest' AND NOT canonical
577654
`
578655
if (manifestChanged && resetCanonicalHistoryOnManifestChange && existing?.['indexed_block'] !== null && existing?.['indexed_block'] !== undefined) {
579-
for (const table of [
580-
'blocks',
581-
'transactions',
582-
'logs',
583-
'contract_discoveries',
584-
'questions',
585-
'pools',
586-
'pool_snapshots',
587-
'pool_state_events',
588-
'vault_snapshots',
589-
'universe_events',
590-
'amm_markets',
591-
'amm_price_snapshots',
592-
'rep_eth_price_snapshots',
593-
'uniswap_rep_eth_markets',
594-
'uniswap_rep_eth_price_observations',
595-
'protocol_timeline_entries',
596-
'open_oracle_report_events',
597-
'escalation_game_events',
598-
'truth_auction_events',
599-
'amm_trade_events',
600-
'fork_migration_events',
601-
'liquidation_approval_events',
602-
'address_activity',
603-
'address_balance_snapshots',
604-
'token_metadata',
605-
])
606-
await transaction.unsafe(`UPDATE ${table} SET canonical = false WHERE chain_id = $1 AND canonical`, [network.chainId])
607-
await transaction`UPDATE entity_state_snapshots SET read_status = 'stale', canonical = false WHERE chain_id = ${network.chainId} AND canonical`
608-
await transaction`UPDATE blocks SET finalized = false WHERE chain_id = ${network.chainId}`
609-
await transaction`UPDATE logs SET finalized = false WHERE chain_id = ${network.chainId}`
610-
await transaction`DELETE FROM log_scan_cursors WHERE chain_id = ${network.chainId}`
611-
await transaction`
612-
UPDATE contracts SET deployment_block = NULL, deployment_timestamp = NULL,
613-
deployment_block_exact = NULL, deployment_checked_block = NULL
614-
WHERE chain_id = ${network.chainId}
615-
`
616-
await transaction`UPDATE contracts SET canonical = false WHERE chain_id = ${network.chainId} AND provenance <> 'manifest'`
656+
await invalidateCanonicalHistory(transaction, network.chainId)
617657
const previousBlock = BigInt(String(existing['indexed_block']))
618658
await transaction`
619659
UPDATE networks SET indexed_block = NULL, indexed_hash = NULL, indexed_timestamp = NULL, finalized_block = NULL, phase = 'backfilling',
@@ -640,8 +680,8 @@ export class ScannerDatabase {
640680
label = CASE WHEN EXCLUDED.provenance = 'manifest' THEN EXCLUDED.label WHEN contracts.provenance = 'manifest' THEN contracts.label ELSE EXCLUDED.label END,
641681
kind = CASE WHEN EXCLUDED.provenance = 'manifest' THEN EXCLUDED.kind WHEN contracts.provenance = 'manifest' THEN contracts.kind ELSE EXCLUDED.kind END,
642682
provenance = CASE WHEN EXCLUDED.provenance = 'manifest' OR contracts.provenance = 'manifest' THEN 'manifest' ELSE EXCLUDED.provenance END,
643-
discovery_block = CASE WHEN EXCLUDED.provenance = 'manifest' THEN NULL WHEN contracts.provenance = 'manifest' THEN contracts.discovery_block ELSE EXCLUDED.discovery_block END,
644-
discovery_tx_hash = CASE WHEN EXCLUDED.provenance = 'manifest' THEN NULL WHEN contracts.provenance = 'manifest' THEN contracts.discovery_tx_hash ELSE EXCLUDED.discovery_tx_hash END,
683+
discovery_block = CASE WHEN EXCLUDED.provenance = 'manifest' AND (contracts.canonical OR contracts.provenance = 'manifest') THEN contracts.discovery_block WHEN EXCLUDED.provenance = 'manifest' THEN NULL WHEN contracts.provenance = 'manifest' THEN contracts.discovery_block ELSE EXCLUDED.discovery_block END,
684+
discovery_tx_hash = CASE WHEN EXCLUDED.provenance = 'manifest' AND (contracts.canonical OR contracts.provenance = 'manifest') THEN contracts.discovery_tx_hash WHEN EXCLUDED.provenance = 'manifest' THEN NULL WHEN contracts.provenance = 'manifest' THEN contracts.discovery_tx_hash ELSE EXCLUDED.discovery_tx_hash END,
645685
canonical = true
646686
`
647687
}
@@ -969,7 +1009,12 @@ export class ScannerDatabase {
9691009
WHERE chain_id = ${chainId} AND (deployment_block > ${ancestor.toString()} OR deployment_checked_block > ${ancestor.toString()})
9701010
`
9711011
await transaction`UPDATE contracts SET deployment_checked_block = NULL WHERE chain_id = ${chainId} AND deployment_checked_block > ${ancestor.toString()}`
972-
await transaction`UPDATE contracts SET canonical = false WHERE chain_id = ${chainId} AND provenance <> 'manifest' AND discovery_block > ${ancestor.toString()}`
1012+
await transaction`
1013+
UPDATE contracts SET canonical = false
1014+
WHERE chain_id = ${chainId} AND provenance <> 'manifest'
1015+
AND (discovery_block IS NULL OR discovery_block >= ${String(checkpoint['start_block'])})
1016+
AND (discovery_block IS NULL OR discovery_block > ${ancestor.toString()})
1017+
`
9731018
await transaction`
9741019
UPDATE contracts AS contract SET
9751020
label = discovery.label,
@@ -1300,6 +1345,33 @@ export class ScannerDatabase {
13001345
})
13011346
}
13021347

1348+
async advanceNetworkStartBlock(chainId: number, startBlock: bigint, lease: IndexerLease): Promise<boolean> {
1349+
return await withIndexerLease(lease, async (transaction) => {
1350+
const rows = await transaction`SELECT start_block, indexed_block FROM networks WHERE chain_id = ${chainId} FOR UPDATE`
1351+
const row = rows[0]
1352+
if (row === undefined) throw new DatabaseConsistencyError(`Network ${chainId} is not initialized`)
1353+
const storedStartBlock = BigInt(String(row['start_block']))
1354+
if (startBlock <= storedStartBlock) return false
1355+
const previousBlock = row['indexed_block'] === null || row['indexed_block'] === undefined ? undefined : BigInt(String(row['indexed_block']))
1356+
const invalidatedDepth = previousBlock === undefined ? 0n : previousBlock - storedStartBlock + 1n
1357+
await invalidateCanonicalHistory(transaction, chainId, startBlock)
1358+
await transaction`
1359+
UPDATE networks SET start_block = ${startBlock.toString()}, indexed_block = NULL, indexed_hash = NULL,
1360+
indexed_timestamp = NULL, finalized_block = NULL, phase = 'backfilling', last_poll_at = now(),
1361+
last_success_at = now(), last_error = NULL, failure_started_at = NULL, consecutive_failures = 0,
1362+
last_reorg_at = now(), last_reorg_depth = ${invalidatedDepth.toString()},
1363+
next_retry_at = NULL, updated_at = now()
1364+
WHERE chain_id = ${chainId}
1365+
`
1366+
await lockLiveEventWriter(transaction)
1367+
await transaction`
1368+
INSERT INTO live_events (event, payload)
1369+
VALUES ('reorg', (${JSON.stringify({ chainId, previousBlock: previousBlock?.toString(), ancestor: '-1', depth: invalidatedDepth.toString(), startBlock: startBlock.toString() })}::text)::jsonb)
1370+
`
1371+
return true
1372+
})
1373+
}
1374+
13031375
async recordFailure(chainId: number, message: string, nextRetryAt: Date, lease: IndexerLease): Promise<void> {
13041376
await withIndexerLease(lease, async (transaction) => {
13051377
const rows = await transaction`

augurScan/src/indexer-runtime.ts

Lines changed: 111 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,18 @@
11
import { type AddressActivity, DatabaseConsistencyError, databaseConsistencyDiagnosticMessage, type IndexerLease, type StoredTransaction } from './database.ts'
22
import { errorChainIncludes } from './error-chain.ts'
3-
import { type Address, type Hash, type Log, type PublicClient, type TransactionReceipt, zeroAddress } from './ethereum.ts'
3+
import {
4+
type Address,
5+
createPublicClient,
6+
type Hash,
7+
http,
8+
type Log,
9+
type PublicClient,
10+
type RpcFetchFn,
11+
type TransactionReceipt,
12+
zeroAddress,
13+
} from './ethereum.ts'
414
import { jsonRpcErrorName, safeRpcProviderMessage } from './logging.ts'
5-
import { RpcRequestMethodError, rpcQueueSaturationFrom } from './rpc-request-queue.ts'
15+
import { RpcRequestMethodError, type RpcRequestQueue, rpcQueueSaturationFrom, withRpcRequestQueue } from './rpc-request-queue.ts'
616
import { bigintToSafeNumber } from './time.ts'
717
import type { ContractMetadata, StoredLog } from './types.ts'
818

@@ -158,6 +168,24 @@ export const isPermanentHistoricalCodeError = (error: unknown): boolean => {
158168
return false
159169
}
160170

171+
export const isPermanentHistoricalLogError = (error: unknown): boolean => {
172+
const seen = new Set<unknown>()
173+
let getLogsRequest = false
174+
let prunedHistory = false
175+
let current: unknown = error
176+
while (typeof current === 'object' && current !== null && !seen.has(current)) {
177+
seen.add(current)
178+
if (current instanceof RpcRequestMethodError && current.method === 'eth_getLogs') getLogsRequest = true
179+
if ('code' in current && current.code === 4444) prunedHistory = true
180+
for (const description of preferredRpcDescriptions(current)) {
181+
const normalized = classifiedRpcDescription(description)
182+
if (normalized.includes('pruned history unavailable') || normalized.includes('historical logs unavailable')) prunedHistory = true
183+
}
184+
current = 'cause' in current ? current.cause : undefined
185+
}
186+
return getLogsRequest && prunedHistory
187+
}
188+
161189
const isPrunedHistoricalStateError = (error: unknown): boolean => {
162190
const seen = new Set<unknown>()
163191
let current: unknown = error
@@ -179,10 +207,74 @@ const isPrunedHistoricalStateError = (error: unknown): boolean => {
179207
}
180208

181209
export const isSplittableLogRangeError = (error: unknown): boolean => {
210+
if (isPermanentHistoricalLogError(error)) return false
182211
const category = rpcErrorCategory(error)
183212
return category !== undefined && category !== 'rate-limit'
184213
}
185214

215+
export const createLogClient = (rpcUrl: string, endpoint: string, queue: RpcRequestQueue, fetchFn: RpcFetchFn, retryDelay?: number): PublicClient =>
216+
createPublicClient({
217+
transport: withRpcRequestQueue(
218+
http(rpcUrl, {
219+
fetchFn,
220+
requestTimeout: 20_000,
221+
retryCount: 2,
222+
...(retryDelay === undefined ? {} : { retryDelay }),
223+
}),
224+
queue,
225+
endpoint,
226+
),
227+
})
228+
229+
export const findEarliestAvailableLogBlock = async (
230+
startBlock: bigint,
231+
observedHead: bigint,
232+
logsAt: (blockNumber: bigint) => Promise<void>,
233+
startBlockKnownUnavailable = false,
234+
): Promise<bigint> => {
235+
if (startBlock > observedHead) throw new Error('The log availability search start must not exceed the observed head')
236+
const isAvailable = async (blockNumber: bigint): Promise<boolean> => {
237+
try {
238+
await logsAt(blockNumber)
239+
return true
240+
} catch (error) {
241+
if (isPermanentHistoricalLogError(error)) return false
242+
throw error
243+
}
244+
}
245+
if (!startBlockKnownUnavailable && (await isAvailable(startBlock))) return startBlock
246+
if (!(await isAvailable(observedHead))) throw new ChainConfigurationError(`RPC cannot serve logs at observed head #${observedHead}`)
247+
let lower = startBlock
248+
let upper = observedHead
249+
while (lower + 1n < upper) {
250+
const middle = lower + (upper - lower) / 2n
251+
if (await isAvailable(middle)) upper = middle
252+
else lower = middle
253+
}
254+
return upper
255+
}
256+
257+
export const findEarliestAvailableLogProvider = async <TProvider>(
258+
providers: readonly TProvider[],
259+
startBlock: bigint,
260+
observedHead: (provider: TProvider) => Promise<bigint>,
261+
logsAt: (provider: TProvider, blockNumber: bigint) => Promise<void>,
262+
): Promise<{ readonly provider: TProvider; readonly startBlock: bigint } | undefined> => {
263+
let earliest: { readonly provider: TProvider; readonly startBlock: bigint } | undefined
264+
for (const provider of providers) {
265+
try {
266+
const head = await observedHead(provider)
267+
if (head < startBlock) continue
268+
const availableStart = await findEarliestAvailableLogBlock(startBlock, head, (blockNumber) => logsAt(provider, blockNumber))
269+
if (earliest === undefined || availableStart < earliest.startBlock) earliest = { provider, startBlock: availableStart }
270+
} catch {
271+
// Recovery is best-effort across providers. The lifecycle retains the
272+
// original failure when none can establish a usable log boundary.
273+
}
274+
}
275+
return earliest
276+
}
277+
186278
export const labelsFrom = (contracts: ReadonlyMap<string, ContractMetadata>): Map<string, string> =>
187279
new Map([['0x0000000000000000000000000000000000000000', 'Zero address'], ...[...contracts].map(([address, contract]) => [address, contract.label] as const)])
188280

@@ -289,6 +381,7 @@ export const withVerifiedProvider = async <TProvider extends ChainProvider, TRes
289381
stopFailover = (_error: unknown): boolean => false,
290382
onAttempt = (_provider: TProvider): void => {},
291383
verifiedProviders?: WeakSet<TProvider>,
384+
onFailure = (_provider: TProvider, _error: unknown): void => {},
292385
): Promise<TResult> => {
293386
let lastFailure: unknown
294387
for (const provider of providers) {
@@ -301,6 +394,7 @@ export const withVerifiedProvider = async <TProvider extends ChainProvider, TRes
301394
}
302395
return await operation(provider)
303396
} catch (error) {
397+
onFailure(provider, error)
304398
if (stopFailover(error)) throw error
305399
lastFailure = error
306400
}
@@ -627,6 +721,7 @@ type NetworkLifecycle = {
627721
readonly verify: () => Promise<void>
628722
readonly poll: () => Promise<boolean>
629723
readonly failure: (message: string, nextRetryAt: Date, reason: string) => Promise<void>
724+
readonly recover?: (error: unknown) => Promise<boolean>
630725
readonly intervalMs: number
631726
readonly signal: AbortSignal
632727
readonly random?: () => number
@@ -650,7 +745,7 @@ export const retryDelayMs = (consecutiveFailures: number, intervalMs: number, ra
650745
return Math.min(Math.round(base * (0.8 + random() * 0.4)), 300_000)
651746
}
652747

653-
export const runNetworkLifecycle = async ({ verify, poll, failure, intervalMs, signal, random, shouldRethrow }: NetworkLifecycle): Promise<void> => {
748+
export const runNetworkLifecycle = async ({ verify, poll, failure, recover, intervalMs, signal, random, shouldRethrow }: NetworkLifecycle): Promise<void> => {
654749
let verified = false
655750
let consecutiveFailures = 0
656751
while (!signal.aborted) {
@@ -666,14 +761,19 @@ export const runNetworkLifecycle = async ({ verify, poll, failure, intervalMs, s
666761
consecutiveFailures = 0
667762
} catch (error) {
668763
if (error instanceof LeaseLostError || shouldRethrow?.(error) === true) throw error
669-
consecutiveFailures++
670-
delayAfterFailure = retryDelayMs(consecutiveFailures, intervalMs, random)
671-
try {
672-
const failureMessage = safeIndexerFailure(error)
673-
const failureReason = failureMessage === 'RPC request failed; retrying' ? rpcIndexerFailureReason(error) : safeIndexerFailureReason(error)
674-
await failure(failureMessage, new Date(Date.now() + delayAfterFailure), failureReason)
675-
} catch (failureError) {
676-
throw new IndexerOwnershipStageError('record-failure', failureError)
764+
if ((await recover?.(error)) === true) {
765+
consecutiveFailures = 0
766+
caughtUp = false
767+
} else {
768+
consecutiveFailures++
769+
delayAfterFailure = retryDelayMs(consecutiveFailures, intervalMs, random)
770+
try {
771+
const failureMessage = safeIndexerFailure(error)
772+
const failureReason = failureMessage === 'RPC request failed; retrying' ? rpcIndexerFailureReason(error) : safeIndexerFailureReason(error)
773+
await failure(failureMessage, new Date(Date.now() + delayAfterFailure), failureReason)
774+
} catch (failureError) {
775+
throw new IndexerOwnershipStageError('record-failure', failureError)
776+
}
677777
}
678778
}
679779
await waitForIndexerDelay(delayAfterFailure ?? (caughtUp ? Math.max(0, intervalMs - (Date.now() - startedAt)) : 0), signal)

0 commit comments

Comments
 (0)