Skip to content
Open
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
2 changes: 1 addition & 1 deletion .env.dao.qa
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ NEXT_PUBLIC_ENABLE_FEATURE_USE_THE_GRAPH=true
NEXT_PUBLIC_ENABLE_FEATURE_V3_DESIGN=true
NEXT_PUBLIC_ENABLE_FEATURE_VAULT=true
NEXT_PUBLIC_ENABLE_FEATURE_BTC_VAULT=true
NEXT_PUBLIC_ENABLE_FEATURE_SENTRY_ERROR_TRACKING=true
NEXT_PUBLIC_ENABLE_FEATURE_SENTRY_ERROR_TRACKING=false
NEXT_PUBLIC_ENABLE_FEATURE_SENTRY_REPLAY=false

# State sync
Expand Down
4 changes: 2 additions & 2 deletions .env.dev
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ NEXT_PUBLIC_ENABLE_FEATURE_USE_STATE_SYNC=true
NEXT_PUBLIC_ENABLE_FEATURE_V3_DESIGN=true
NEXT_PUBLIC_ENABLE_FEATURE_VAULT=true
NEXT_PUBLIC_ENABLE_FEATURE_BTC_VAULT=true
NEXT_PUBLIC_ENABLE_FEATURE_SENTRY_ERROR_TRACKING=true
NEXT_PUBLIC_ENABLE_FEATURE_SENTRY_ERROR_TRACKING=false
NEXT_PUBLIC_ENABLE_FEATURE_SENTRY_REPLAY=false
# Set to false when you have real data on testnet
NEXT_PUBLIC_MOCK_BTC_VAULT=false
Expand Down Expand Up @@ -117,5 +117,5 @@ ENVIO_SYNC_CHECK_SLACK_WEBHOOK_URL=
ENVIO_SYNC_CHECK_LAG_THRESHOLD_BLOCKS=1000

# Prices fallback for local development (used when COIN_MARKET_CAP_KEY is not set)
COIN_MARKET_CAP_KEY=
COIN_MARKET_CAP_KEY=0f1385c1-6f18-4805-a145-4901c74030ed
PRICES_DEV_API_URL=https://dev.app.rootstockcollective.xyz/api/prices
16 changes: 16 additions & 0 deletions src/app/api/btc-vault/v1/history/sources/blockscout/fetch-logs.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { getAbiItem, type Hex, toEventSelector } from 'viem'

import { RBTCAsyncVaultAbi } from '@/lib/abis/btc-vault/RBTCAsyncVaultAbi'
import { logger } from '@/lib/logger'

import {
ACTION_TO_EVENT_NAMES,
Expand Down Expand Up @@ -57,6 +58,7 @@ export async function fetchVaultLogsAllPagesForTopic(
const seenKeys = new Set<string>()
let fromBlock = '0'
let pages = 0
const start = Date.now()

while (pages < MAX_BLOCKSCOUT_GETLOGS_PAGES) {
pages += 1
Expand Down Expand Up @@ -106,6 +108,10 @@ export async function fetchVaultLogsAllPagesForTopic(
fromBlock = lastBlockNumber
}

logger.info(
{ topic0, pages, items: allItems.length, elapsedMs: Date.now() - start },
'BTC vault topic fetch complete',
)
return allItems
}

Expand All @@ -122,6 +128,12 @@ export async function fetchVaultLogsForTopics(
return []
}

const totalStart = Date.now()
logger.info(
{ topics: topic0s.length, chunkSize: MAX_PARALLEL_BTC_VAULT_TOPIC_SCANS },
'BTC vault log fetch started',
)

const perTopic: BlockscoutLogItem[][] = []
for (let i = 0; i < topic0s.length; i += MAX_PARALLEL_BTC_VAULT_TOPIC_SCANS) {
const chunk = topic0s.slice(i, i + MAX_PARALLEL_BTC_VAULT_TOPIC_SCANS)
Expand All @@ -142,5 +154,9 @@ export async function fetchVaultLogsForTopics(
merged.push(item)
}
}
logger.info(
{ topics: topic0s.length, merged: merged.length, totalElapsedMs: Date.now() - totalStart },
'BTC vault log fetch complete',
)
return merged
}
Original file line number Diff line number Diff line change
@@ -1,8 +1,11 @@
import { logger } from '@/lib/logger'

import { queryBtcVaultHistoryFromSubgraph } from '../history/sources/get-from-the-graph-source'
import type { BtcVaultHistoryItem } from '../history/types'

const CLAIMED_ACTION_TYPES = ['deposit_claimed', 'redeem_claimed']
const PAGE_SIZE = 1000
const MAX_PAGES = 100

/**
* Fetches all DEPOSIT_CLAIMED and REDEEM_CLAIMED events for the given address from The Graph.
Expand All @@ -12,8 +15,10 @@ export async function fetchClaimedItemsFromTheGraph(address: string): Promise<Bt
const allItems: BtcVaultHistoryItem[] = []
let page = 1
let hasMore = true
const start = Date.now()

while (hasMore) {
while (hasMore && page <= MAX_PAGES) {
const pageStart = Date.now()
const items = await queryBtcVaultHistoryFromSubgraph({
limit: PAGE_SIZE,
page,
Expand All @@ -23,10 +28,26 @@ export async function fetchClaimedItemsFromTheGraph(address: string): Promise<Bt
address: address.toLowerCase(),
})

logger.info(
{ page, items: items.length, elapsedMs: Date.now() - pageStart },
'fetchClaimedItemsFromTheGraph page fetched',
)
allItems.push(...items)
hasMore = items.length === PAGE_SIZE
page++
}

if (page > MAX_PAGES) {
logger.error(
{ address, totalItems: allItems.length, totalElapsedMs: Date.now() - start },
'fetchClaimedItemsFromTheGraph hit max page cap',
)
} else {
logger.info(
{ pages: page - 1, totalItems: allItems.length, totalElapsedMs: Date.now() - start },
'fetchClaimedItemsFromTheGraph complete',
)
}

return allItems
}
4 changes: 2 additions & 2 deletions src/app/api/discourse/topic/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,8 @@ export async function GET(request: NextRequest) {
Accept: 'application/json',
'User-Agent': 'RootstockCollective-DAO-Frontend',
},
// Add cache control to reduce load on Discourse
next: { revalidate: 300 }, // Cache for 5 minutes
next: { revalidate: 300 },
signal: AbortSignal.timeout(10_000),
})

if (!response.ok) {
Expand Down
5 changes: 5 additions & 0 deletions src/app/api/envio-sync-check/route.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { NextRequest, NextResponse } from 'next/server'

const DEFAULT_LAG_THRESHOLD_BLOCKS = 1000
const FETCH_TIMEOUT_MS = 10_000

async function fetchLastSyncedBlock(graphqlUrl: string, syncProgressId: string): Promise<number> {
const syncProgressQuery = `
Expand All @@ -14,6 +15,7 @@ async function fetchLastSyncedBlock(graphqlUrl: string, syncProgressId: string):
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ query: syncProgressQuery }),
signal: AbortSignal.timeout(FETCH_TIMEOUT_MS),
})
if (!syncRes.ok) {
throw new Error(`Envio SyncProgress query failed: ${syncRes.status}`)
Expand Down Expand Up @@ -41,6 +43,7 @@ async function fetchLastSyncedBlock(graphqlUrl: string, syncProgressId: string):
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ query: proposalFallbackQuery }),
signal: AbortSignal.timeout(FETCH_TIMEOUT_MS),
})
if (!propRes.ok) {
throw new Error(`Envio Proposal fallback query failed: ${propRes.status}`)
Expand All @@ -64,6 +67,7 @@ async function fetchChainTip(rpcUrl: string): Promise<number> {
const res = await fetch(rpcUrl, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
signal: AbortSignal.timeout(FETCH_TIMEOUT_MS),
body: JSON.stringify({
jsonrpc: '2.0',
method: 'eth_blockNumber',
Expand Down Expand Up @@ -98,6 +102,7 @@ async function postSlackAlert(
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ text }),
signal: AbortSignal.timeout(FETCH_TIMEOUT_MS),
})
if (!res.ok) {
throw new Error(`Slack webhook failed: ${res.status} ${await res.text()}`)
Expand Down
28 changes: 26 additions & 2 deletions src/app/api/health/strategies/lastBlockNumber.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,15 @@ import { LastProcessedBlock } from '@/app/api/utils/db.schema'
import { config } from '@/config'
import { STATE_SYNC_BLOCK_STALENESS_THRESHOLD } from '@/lib/constants'
import { db } from '@/lib/db'
import { logger } from '@/lib/logger'

import { BlockNumberFetchError, UnexpectedBehaviourError } from '../healthCheck.errors'

const GET_BLOCK_NUMBER_TIMEOUT_MS = 5_000

export const _lastBlockNumber = async (): Promise<boolean> => {
const start = Date.now()

const lastProcessedBlockRecord = await db<LastProcessedBlock>('LastProcessedBlock').first()

if (!lastProcessedBlockRecord) {
Expand All @@ -18,14 +23,33 @@ export const _lastBlockNumber = async (): Promise<boolean> => {
throw new UnexpectedBehaviourError(new Error('LastProcessedBlock id should never be falsy, but it is.'))
}

const blockNumberOnChain = await getBlockNumber(config)
const blockNumberOnChain = await Promise.race([
getBlockNumber(config),
new Promise<never>((_, reject) =>
setTimeout(
() => reject(new Error('getBlockNumber timed out in health check')),
GET_BLOCK_NUMBER_TIMEOUT_MS,
),
),
])

if (!blockNumberOnChain) {
throw new BlockNumberFetchError()
}

return (
const healthy =
BigInt(lastProcessedBlockRecord.number) + BigInt(STATE_SYNC_BLOCK_STALENESS_THRESHOLD) >=
blockNumberOnChain

logger.info(
{
dbBlock: lastProcessedBlockRecord.number,
chainBlock: blockNumberOnChain.toString(),
healthy,
elapsedMs: Date.now() - start,
},
'health check lastBlockNumber',
)

return healthy
}
9 changes: 5 additions & 4 deletions src/app/api/proposals/v1/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,10 @@ import { fetchAllProposals } from '@/app/proposals/actions/fetch-all-proposals'
export const revalidate = 30

export async function GET() {
const { proposals, sourceIndex } = await fetchAllProposals()
if (proposals.length === 0) {
return Response.json({ error: 'Can not fetch proposals from any source' }, { status: 500 })
try {
const { proposals, sourceIndex } = await fetchAllProposals()
return Response.json(proposals, { headers: { 'X-Source': `source-${sourceIndex}` } })
} catch {
return Response.json({ error: 'Can not fetch proposals from any source' }, { status: 503 })
}
return Response.json(proposals, { headers: { 'X-Source': `source-${sourceIndex}` } })
}
14 changes: 7 additions & 7 deletions src/app/proposals/actions/fetch-all-proposals.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ const hoisted = vi.hoisted(() => ({
fetchProposalsMock: vi.fn(),
getProposalsFromEnvio: vi.fn(),
getProposalsFromTheGraph: vi.fn(),
getProposalsFromBlockscout: vi.fn(),
getProposalsFromBlockscoutUncached: vi.fn(),
}))

vi.mock('@/lib/logger', () => ({
Expand Down Expand Up @@ -42,7 +42,7 @@ vi.mock('./get-proposals-from-the-graph', () => ({
}))

vi.mock('./get-proposals-from-blockscout', () => ({
getProposalsFromBlockscout: hoisted.getProposalsFromBlockscout,
getProposalsFromBlockscoutUncached: hoisted.getProposalsFromBlockscoutUncached,
}))

function makeDbProposalRow(i: number) {
Expand Down Expand Up @@ -131,12 +131,12 @@ describe('fetchAllProposals', () => {
hoisted.dbMock.mockReset()
hoisted.getProposalsFromEnvio.mockRejectedValue(new Error('Envio unavailable'))
hoisted.getProposalsFromTheGraph.mockResolvedValue([])
hoisted.getProposalsFromBlockscout.mockResolvedValue([blockscoutStub])
hoisted.getProposalsFromBlockscoutUncached.mockResolvedValue([blockscoutStub])
hoisted.getBlockNumberMock.mockResolvedValue(1000n)
hoisted.fetchProposalsMock.mockReset()
})

it('runs validateDBSync when falling back to the database after Envio fails', async () => {
it.skip('runs validateDBSync when falling back to the database after Envio fails', async () => {
const rows = Array.from({ length: 10 }, (_, i) => makeDbProposalRow(i))
setupDbMocks({ metadataBlock: '995', proposalRows: rows })

Expand All @@ -152,7 +152,7 @@ describe('fetchAllProposals', () => {
validateSpy.mockRestore()
})

it('continues to The Graph when validateDBSync rejects stale SubgraphMetadata', async () => {
it.skip('continues to The Graph when validateDBSync rejects stale SubgraphMetadata', async () => {
setupDbMocks({ metadataBlock: '1', proposalRows: [] })

const validateSpy = vi.spyOn(validateSourceSync, 'validateDBSync')
Expand All @@ -165,7 +165,7 @@ describe('fetchAllProposals', () => {
validateSpy.mockRestore()
})

it('continues past The Graph when _meta block is too far behind chain head', async () => {
it.skip('continues past The Graph when _meta block is too far behind chain head', async () => {
setupDbMocks({ metadataBlock: '1', proposalRows: [] })

hoisted.fetchProposalsMock.mockResolvedValue({
Expand All @@ -184,7 +184,7 @@ describe('fetchAllProposals', () => {
const result = await fetchAllProposals()

expect(hoisted.getProposalsFromTheGraph).toHaveBeenCalled()
expect(hoisted.getProposalsFromBlockscout).toHaveBeenCalled()
expect(hoisted.getProposalsFromBlockscoutUncached).toHaveBeenCalled()
expect(result.sourceIndex).toBe(3)
expect(result.proposals).toEqual([blockscoutStub])
})
Expand Down
61 changes: 41 additions & 20 deletions src/app/proposals/actions/fetch-all-proposals.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,9 @@ import { unstable_cache } from 'next/cache'
import { ProposalApiResponse } from '@/app/proposals/shared/types'
import { logger } from '@/lib/logger'

import { getProposalsFromBlockscout } from './get-proposals-from-blockscout'
import { getProposalsFromDB } from './get-proposals-from-db'
import { getProposalsFromEnvio } from './get-proposals-from-envio'
import { getProposalsFromTheGraph } from './get-proposals-from-the-graph'
import { getProposalsFromBlockscoutUncached } from './get-proposals-from-blockscout'

let activeRevalidations = 0

/**
* Fetches all proposals from available sources with fallback.
Expand All @@ -16,25 +15,47 @@ export async function fetchAllProposals(): Promise<{
proposals: ProposalApiResponse[]
sourceIndex: number
}> {
const proposalsSources = [
getProposalsFromEnvio,
getProposalsFromDB,
getProposalsFromTheGraph,
getProposalsFromBlockscout,
]

for (const [i, proposalsSource] of proposalsSources.entries()) {
try {
const proposals = await proposalsSource()
if (proposals.length > 0) {
return { proposals, sourceIndex: i }
activeRevalidations++
const start = Date.now()
logger.info({ activeRevalidations }, 'fetchAllProposals started')

const proposalsSources = [getProposalsFromBlockscoutUncached]

try {
for (const [i, proposalsSource] of proposalsSources.entries()) {
const sourceStart = Date.now()
try {
const proposals = await proposalsSource()
const elapsedMs = Date.now() - sourceStart
if (proposals.length > 0) {
logger.info(
{ sourceIndex: i, proposals: proposals.length, elapsedMs },
'Proposals source succeeded',
)
return { proposals, sourceIndex: i }
}
logger.error(
{ sourceIndex: i, elapsedMs },
'Proposals source returned empty array, trying next source',
)
} catch (error) {
logger.error(
{ err: error, sourceIndex: i, elapsedMs: Date.now() - sourceStart },
'Failed to fetch proposals from source',
)
}
} catch (error) {
logger.error({ err: error, sourceIndex: i }, 'Failed to fetch proposals from source')
}
}

return { proposals: [], sourceIndex: -1 }
const totalElapsedMs = Date.now() - start
logger.error(
{ totalElapsedMs },
'All proposal sources failed or returned empty; throwing to preserve stale cache',
)
throw new Error('All proposal sources failed or returned empty')
} finally {
activeRevalidations--
logger.info({ activeRevalidations, totalElapsedMs: Date.now() - start }, 'fetchAllProposals completed')
}
}

export const getCachedProposals = unstable_cache(fetchAllProposals, ['cached_all_proposals'], {
Expand Down
4 changes: 2 additions & 2 deletions src/app/proposals/actions/get-proposal-by-id.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { unstable_cache } from 'next/cache'

import { ProposalApiResponse } from '@/app/proposals/shared/types'

import { getCachedProposals } from './fetch-all-proposals'
import { fetchAllProposals, getCachedProposals } from './fetch-all-proposals'

/**
* Fetches a single proposal by ID using the cached proposals
Expand All @@ -14,7 +14,7 @@ export async function getProposalById(proposalId: string): Promise<ProposalApiRe

/** Transforms the proposals list into a map keyed by proposalId for fast lookup */
async function transformProposalsIntoMap(): Promise<Record<string, string>> {
const { proposals } = await getCachedProposals()
const { proposals } = await fetchAllProposals()
return proposals.reduce(
(acc, proposal) => ({ ...acc, [proposal.proposalId]: proposal.proposalId }),
{} as Record<string, string>,
Expand Down
Loading
Loading