Skip to content
Merged
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
37 changes: 24 additions & 13 deletions apps/indexer/src/main/normalized_replay_catchup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ use crate::{
HeaderAuditMode, RawFactNormalizedEventReplayOutcome, RawFactNormalizedEventReplayRequest,
RawFactNormalizedEventReplaySelection, active_closure_or_dependency_replay_adapters,
chain_has_closure_or_dependency_replay_adapter, replay_raw_fact_normalized_events,
sync_full_closure_normalized_events_from_persisted_raw_payloads,
select_log_bounded_replay_to_block,
sync_automatic_two_phase_full_closure_normalized_events,
unsupported_closure_replay_adapters,
},
};
Expand All @@ -35,7 +36,10 @@ mod sources;
mod test_hook;

#[cfg(test)]
pub(crate) use test_hook::install_after_rewind as install_after_rewind_test_hook;
pub(crate) use test_hook::{
install_after_coverage_recovery as install_after_coverage_recovery_test_hook,
install_after_rewind as install_after_rewind_test_hook,
};

use coverage_recovery::replay_full_closure_with_coverage_recovery;
use cursors::{
Expand All @@ -46,7 +50,7 @@ use indexes::{
ensure_projection_indexes_after_catchup, prepare_deferred_projection_indexes_for_fresh_replay,
restore_deferred_projection_indexes,
};
use sources::{load_canonical_raw_log_bounds, select_log_bounded_replay_to_block};
use sources::load_canonical_raw_log_bounds;

pub(crate) const DEFAULT_NORMALIZED_REPLAY_CATCHUP_CHUNK_BLOCKS: i64 = 262_144;
pub(crate) const DEFAULT_NORMALIZED_REPLAY_CATCHUP_MAX_LOGS_PER_CHUNK: usize = 100_000;
Expand Down Expand Up @@ -520,6 +524,7 @@ async fn replay_full_closure_or_dependency_normalized_events(
chain: &str,
from_block: i64,
to_block: i64,
stateless_ranges: &[(i64, i64)],
max_raw_logs_per_page: usize,
) -> Result<RawFactNormalizedEventReplayOutcome> {
let adapters = active_closure_or_dependency_replay_adapters(pool, chain).await?;
Expand All @@ -536,35 +541,41 @@ async fn replay_full_closure_or_dependency_normalized_events(
chain,
from_block,
to_block,
stateless_range_count = stateless_ranges.len(),
stateless_ranges = ?stateless_ranges,
max_raw_logs_per_page,
adapter_count = adapters.len(),
adapters = ?adapters,
"full closure normalized-event replay session started"
"two-phase full closure normalized-event replay session started"
);

let replay = sync_full_closure_normalized_events_from_persisted_raw_payloads(
let replay = sync_automatic_two_phase_full_closure_normalized_events(
pool,
deployment_profile,
chain,
CURSOR_KIND_RAW_FACT_NORMALIZED_EVENTS,
from_block,
to_block,
stateless_ranges,
&adapters,
max_raw_logs_per_page,
)
.await?;
let summary = replay.summary;
let stateless = replay.stateless;
let closure = replay.closure;

Ok(RawFactNormalizedEventReplayOutcome {
deployment_profile: deployment_profile.to_owned(),
chain: chain.to_owned(),
selection_kind: "full_closure",
selection_kind: "two_phase_full_closure",
source_scope_target_count: adapters.len(),
selected_block_count: 0,
canonical_raw_log_count: summary.scanned_log_count,
scanned_raw_log_count: summary.scanned_log_count,
matched_raw_log_count: summary.matched_log_count,
normalized_event_synced_count: summary.total_synced_count,
normalized_event_inserted_count: summary.total_inserted_count,
selected_block_count: stateless.selected_block_count,
canonical_raw_log_count: stateless.canonical_raw_log_count,
scanned_raw_log_count: stateless.scanned_raw_log_count + closure.scanned_log_count,
matched_raw_log_count: stateless.matched_raw_log_count + closure.matched_log_count,
normalized_event_synced_count: stateless.normalized_event_synced_count
+ closure.total_synced_count,
normalized_event_inserted_count: stateless.normalized_event_inserted_count
+ closure.total_inserted_count,
})
}
128 changes: 126 additions & 2 deletions apps/indexer/src/main/normalized_replay_catchup/coverage_recovery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ use crate::{
provider::ChainProviderOps,
reconciliation::{
EnsV2LiveCoverageRecoveryStatus, HeaderAuditMode, RawFactNormalizedEventReplayOutcome,
recover_ens_v2_live_coverage_requirement,
automatic_stateless_replay_completed, recover_ens_v2_live_coverage_requirement,
},
};

Expand All @@ -31,20 +31,23 @@ pub(super) async fn replay_full_closure_with_coverage_recovery(
RawLogStagingInputVersion,
)> {
let mut recovery_attempt = 0_usize;
let mut stateless_ranges = vec![(from_block, to_block)];
loop {
let replay_error = match replay_full_closure_or_dependency_normalized_events(
pool,
deployment_profile,
chain,
from_block,
to_block,
&stateless_ranges,
max_raw_logs_per_page,
)
.await
{
Ok(outcome) => return Ok((outcome, raw_log_input_version)),
Err(error) => error,
};
let stateless_replay_completed = automatic_stateless_replay_completed(&replay_error);
let Some(requirement) = bigname_adapters::ens_v2_missing_coverage(&replay_error).cloned()
else {
return Err(replay_error);
Expand Down Expand Up @@ -83,8 +86,60 @@ pub(super) async fn replay_full_closure_with_coverage_recovery(
));
}

raw_log_input_version =
// Preserve the original full span when preflight validation prevented
// phase one from running. Once phase one completed, retain only every
// exact span fetched by later recovery attempts. The stateful adapter
// pass still restarts over its complete span.
if stateless_replay_completed {
stateless_ranges.clear();
}
include_stateless_range(
&mut stateless_ranges,
requirement.required_from_block,
requirement.required_to_block,
);
#[cfg(test)]
super::test_hook::pause_after_coverage_recovery(pool, deployment_profile, chain).await;

let observed_raw_log_input_version =
bigname_storage::load_raw_log_staging_input_version(pool, chain).await?;
if observed_raw_log_input_version.retention_generation
!= raw_log_input_version.retention_generation
{
return Err(replay_error.context(format!(
"raw-log retention generation changed during normalized replay coverage recovery: expected {}, observed {}; replan the replay from current authority",
raw_log_input_version.retention_generation,
observed_raw_log_input_version.retention_generation,
)));
}
if from_block > 0
&& raw_log_changed_outside_stateless_ranges(
pool,
chain,
raw_log_input_version.revision,
&stateless_ranges,
0,
from_block - 1,
)
.await?
{
return Err(replay_error.context(format!(
"raw-log staging input changed below normalized replay range start {from_block} during coverage recovery; replan from the durable cursor"
)));
}
let widened_for_concurrent_input = raw_log_changed_outside_stateless_ranges(
pool,
chain,
raw_log_input_version.revision,
&stateless_ranges,
from_block,
to_block,
)
.await?;
if widened_for_concurrent_input {
include_stateless_range(&mut stateless_ranges, from_block, to_block);
}
raw_log_input_version = observed_raw_log_input_version;
info!(
service = "indexer",
command = "run",
Expand All @@ -96,7 +151,76 @@ pub(super) async fn replay_full_closure_with_coverage_recovery(
to_block = requirement.required_to_block,
retention_generation = requirement.retention_generation,
recovery_attempt,
widened_for_concurrent_input,
stateless_range_count = stateless_ranges.len(),
stateless_ranges = ?stateless_ranges,
"retrying unchanged normalized replay after exact generation-bound coverage recovery"
);
}
}

async fn raw_log_changed_outside_stateless_ranges(
pool: &sqlx::PgPool,
chain: &str,
revision: i64,
stateless_ranges: &[(i64, i64)],
inspected_from_block: i64,
inspected_to_block: i64,
) -> Result<bool> {
let mut next_uncovered_block = inspected_from_block;
for &(range_from_block, range_to_block) in stateless_ranges {
if range_to_block < next_uncovered_block {
continue;
}
if range_from_block > inspected_to_block {
break;
}
let covered_from_block = range_from_block.max(inspected_from_block);
if next_uncovered_block < covered_from_block
&& bigname_storage::raw_log_staging_block_range_changed_since(
pool,
chain,
revision,
next_uncovered_block,
covered_from_block - 1,
)
.await?
{
return Ok(true);
}
let Some(after_covered_block) = range_to_block.checked_add(1) else {
return Ok(false);
};
next_uncovered_block = next_uncovered_block.max(after_covered_block);
if next_uncovered_block > inspected_to_block {
return Ok(false);
}
}
bigname_storage::raw_log_staging_block_range_changed_since(
pool,
chain,
revision,
next_uncovered_block,
inspected_to_block,
)
.await
}

fn include_stateless_range(ranges: &mut Vec<(i64, i64)>, from_block: i64, to_block: i64) {
debug_assert!(from_block <= to_block);
ranges.push((from_block, to_block));
ranges.sort_unstable();

let mut merged = Vec::<(i64, i64)>::with_capacity(ranges.len());
for (from_block, to_block) in ranges.drain(..) {
if let Some((_, merged_to_block)) = merged.last_mut()
&& (from_block <= *merged_to_block
|| merged_to_block.checked_add(1) == Some(from_block))
{
*merged_to_block = (*merged_to_block).max(to_block);
} else {
merged.push((from_block, to_block));
}
}
*ranges = merged;
}
72 changes: 0 additions & 72 deletions apps/indexer/src/main/normalized_replay_catchup/sources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -81,75 +81,3 @@ pub(super) async fn load_canonical_raw_log_bounds(
target_block,
}))
}

pub(super) async fn select_log_bounded_replay_to_block(
pool: &PgPool,
chain: &str,
from_block: i64,
hard_to_block: i64,
max_raw_logs_per_chunk: usize,
) -> Result<i64> {
if from_block >= hard_to_block {
return Ok(hard_to_block);
}
let max_raw_logs_per_chunk = i64::try_from(max_raw_logs_per_chunk)
.context("normalized replay max logs per chunk does not fit in i64")?;

sqlx::query_scalar::<_, i64>(
r#"
WITH ordered_logs AS (
SELECT raw_logs.block_number
FROM raw_logs
JOIN chain_lineage AS lineage
ON lineage.chain_id = raw_logs.chain_id
AND lineage.block_hash = raw_logs.block_hash
WHERE raw_logs.chain_id = $1
AND raw_logs.block_number >= $2
AND raw_logs.block_number <= $3
AND lineage.canonicality_state IN (
'canonical'::canonicality_state,
'safe'::canonicality_state,
'finalized'::canonicality_state
)
AND raw_logs.canonicality_state IN (
'canonical'::canonicality_state,
'safe'::canonicality_state,
'finalized'::canonicality_state
)
ORDER BY raw_logs.block_number ASC, raw_logs.log_index ASC
LIMIT ($4 + 1)
),
numbered_logs AS (
SELECT block_number, ROW_NUMBER() OVER () AS ordinal
FROM ordered_logs
),
overflow AS (
SELECT block_number
FROM numbered_logs
WHERE ordinal = $4 + 1
),
bounded AS (
SELECT block_number
FROM numbered_logs
WHERE NOT EXISTS (SELECT 1 FROM overflow)
OR block_number < (SELECT block_number FROM overflow)
UNION ALL
SELECT MIN(block_number)
FROM numbered_logs
)
SELECT COALESCE(MAX(block_number), $3)
FROM bounded
"#,
)
.bind(chain)
.bind(from_block)
.bind(hard_to_block)
.bind(max_raw_logs_per_chunk)
.fetch_one(pool)
.await
.with_context(|| {
format!(
"failed to select log-bounded normalized replay range for chain {chain} range {from_block}..={hard_to_block}"
)
})
}
Loading
Loading