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
20 changes: 19 additions & 1 deletion src/common/consolidations.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import logging

from eth_typing import BlockNumber, HexStr
from sw_utils import ChainHead
from web3 import Web3
Expand All @@ -9,6 +11,8 @@
from src.config.settings import settings
from src.validators.typings import ConsensusValidator

logger = logging.getLogger(__name__)


async def get_pending_consolidations(
chain_head: ChainHead, consensus_validators: list[ConsensusValidator]
Expand All @@ -33,7 +37,21 @@ async def get_pending_consolidations(
continue

if has_source and not has_target:
raise ValueError(f'Target validator index {target_index} not found in vault validators')
# The submission flow guarantees consolidation targets are vault-registered
# validators, so an unresolved target here means `consensus_validators`
# (built from the local SQLite DB) is lagging the on-chain registration event
# -- e.g. after a restart or during event-scan catch-up. Warn instead of
# raising: aborting would stall the whole ValidatorTask.process_block cycle
# (registration + withdrawals) until indexing catches up. Both indexes are
# still known from the CL entry, so keep the source excluded from exit/partial
# selection; the target index is inert since it never matches a real validator.
logger.warning(
'Pending consolidation target validator index %s not found in vault '
'validators (source index %s); vault validator database may be lagging '
'on-chain registration',
target_index,
source_index,
)

result.append(PendingConsolidation(source_index=source_index, target_index=target_index))

Expand Down
150 changes: 150 additions & 0 deletions src/common/tests/test_consolidations.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
from unittest.mock import AsyncMock, patch

import pytest
from sw_utils import ValidatorStatus
from sw_utils.tests import faker

from src.common.consolidations import get_pending_consolidations
from src.common.tests.factories import create_chain_head
from src.common.typings import PendingConsolidation
from src.config.settings import settings
from src.validators.tests.factories import create_consensus_validator


@pytest.mark.usefixtures('fake_settings')
class TestGetPendingConsolidations:
"""CL-queue entries come from `consensus_client.get_pending_consolidations`;
EL-queue entries come from `get_execution_consolidations`. Both are mocked so no
beacon/storage reads happen."""

async def test_cl_queue_unresolved_target_does_not_raise(self):
"""A vault-registered source whose target isn't in `consensus_validators` yet
(local DB indexing race) must not raise -- the entry is kept, with both indexes,
so the source stays excluded from exit/partial selection elsewhere."""
chain_head = create_chain_head()
source_validator = create_consensus_validator(
index=10, status=ValidatorStatus.ACTIVE_ONGOING, activation_epoch=1
)

mock_consensus = AsyncMock()
mock_consensus.get_pending_consolidations.return_value = [
{'source_index': '10', 'target_index': '99'}
]

with patch('src.common.consolidations.consensus_client', mock_consensus), patch(
'src.common.consolidations.get_execution_consolidations',
AsyncMock(return_value=[]),
):
result = await get_pending_consolidations(chain_head, [source_validator])

assert result == [PendingConsolidation(source_index=10, target_index=99)]

async def test_cl_queue_vault_source_and_target_included(self):
chain_head = create_chain_head()
source_validator = create_consensus_validator(
index=10, status=ValidatorStatus.ACTIVE_ONGOING, activation_epoch=1
)
target_validator = create_consensus_validator(
index=20, status=ValidatorStatus.ACTIVE_ONGOING, activation_epoch=1
)

mock_consensus = AsyncMock()
mock_consensus.get_pending_consolidations.return_value = [
{'source_index': '10', 'target_index': '20'}
]

with patch('src.common.consolidations.consensus_client', mock_consensus), patch(
'src.common.consolidations.get_execution_consolidations',
AsyncMock(return_value=[]),
):
result = await get_pending_consolidations(
chain_head, [source_validator, target_validator]
)

assert result == [PendingConsolidation(source_index=10, target_index=20)]

async def test_cl_queue_both_non_vault_skipped(self):
chain_head = create_chain_head()

mock_consensus = AsyncMock()
mock_consensus.get_pending_consolidations.return_value = [
{'source_index': '10', 'target_index': '20'}
]

with patch('src.common.consolidations.consensus_client', mock_consensus), patch(
'src.common.consolidations.get_execution_consolidations',
AsyncMock(return_value=[]),
):
result = await get_pending_consolidations(chain_head, [])

assert result == []

async def test_el_queue_vault_source_address_appended(self):
chain_head = create_chain_head()
source_validator = create_consensus_validator(
index=1,
public_key='0xsource',
status=ValidatorStatus.ACTIVE_ONGOING,
activation_epoch=1,
)
target_validator = create_consensus_validator(
index=2,
public_key='0xtarget',
status=ValidatorStatus.ACTIVE_ONGOING,
activation_epoch=1,
)

mock_consensus = AsyncMock()
mock_consensus.get_pending_consolidations.return_value = []
execution_consolidations = [
{
'source_address': settings.vault,
'source_pubkey': '0xsource',
'target_pubkey': '0xtarget',
}
]

with patch('src.common.consolidations.consensus_client', mock_consensus), patch(
'src.common.consolidations.get_execution_consolidations',
AsyncMock(return_value=execution_consolidations),
):
result = await get_pending_consolidations(
chain_head, [source_validator, target_validator]
)

assert result == [PendingConsolidation(source_index=1, target_index=2)]

async def test_el_queue_wrong_source_address_skipped(self):
chain_head = create_chain_head()
source_validator = create_consensus_validator(
index=1,
public_key='0xsource',
status=ValidatorStatus.ACTIVE_ONGOING,
activation_epoch=1,
)
target_validator = create_consensus_validator(
index=2,
public_key='0xtarget',
status=ValidatorStatus.ACTIVE_ONGOING,
activation_epoch=1,
)

mock_consensus = AsyncMock()
mock_consensus.get_pending_consolidations.return_value = []
execution_consolidations = [
{
'source_address': faker.eth_address(),
'source_pubkey': '0xsource',
'target_pubkey': '0xtarget',
}
]

with patch('src.common.consolidations.consensus_client', mock_consensus), patch(
'src.common.consolidations.get_execution_consolidations',
AsyncMock(return_value=execution_consolidations),
):
result = await get_pending_consolidations(
chain_head, [source_validator, target_validator]
)

assert result == []
Loading