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
75 changes: 67 additions & 8 deletions src/cortex/master.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,15 @@ def validate(self) -> None:
raise ServiceError(503, "invalid chain epoch boundary")


# How far below a completed epoch's boundary the search walks.
#
# The chain retains epoch state for a few epochs, and beyond that answers 0 for
# any block, so a window that fell outside retention cannot be delimited and is
# reported rather than guessed at. This bound is generous relative to a tempo of
# a few hundred blocks and keeps a stuck loop from walking the whole chain.
_END_BLOCK_SEARCH = 8192


class EpochProvider(Protocol):
async def epoch_state(self, netuid: int) -> EpochState: ...

Expand All @@ -91,20 +100,67 @@ def read():
return await asyncio.to_thread(read)

async def end_block(self, netuid: int, epoch: int, start: int, state: EpochState) -> int:
"""Locate the inclusive end even if a restart skipped several epoch transitions."""
"""Locate the inclusive end of a completed epoch.

The window is bracketed by the chain's own boundary — `last_epoch_block`
carries the epoch *after* the one being closed — and by the earliest
block below it that still answers. The remembered `start` is used when
the chain can still read it, and ignored when it cannot: a start
recorded on an earlier tick can point at a block whose state is pruned,
and a search that insists on reading `epoch` there refuses on every tick
and seals nothing.

Measured on netuid 100: the block a journal had remembered for epoch
25267 read as index 0, and the epoch's real end was 200 blocks below
`last_epoch_block`. Nine thousand blocks below, the state is gone.
"""

def index_at(block: int) -> int | None:
"""The epoch index at a block, or None when the chain cannot answer.

A node keeps recent state and prunes the rest, and asking about a
pruned block raises rather than answering. That is a limit of the
read, not a verdict about the window.
"""
try:
index = self.subtensor.get_subnet_epoch_index(netuid, block=block)
except Exception:
return None
return index if isinstance(index, int) else None

def read():
low, high = start, state.last_epoch_block
high = state.last_epoch_block
anchor = self.subtensor.get_block_hash(high)
if self.subtensor.get_subnet_epoch_index(netuid, block=low) != epoch:
raise ServiceError(503, "historical epoch start changed")
if self.subtensor.get_subnet_epoch_index(netuid, block=high) <= epoch:
head = index_at(high)
if head is None or head <= epoch:
raise ServiceError(503, "epoch has not ended")

# Bracket the search. The remembered start is the floor when it is
# still readable and still in this epoch; otherwise the descent is
# bounded, because past the retained window the answer is gone.
low = start if 0 <= start < high else high - 1
if index_at(low) != epoch:
low = high - 1
for _ in range(_END_BLOCK_SEARCH):
if index_at(low) == epoch:
break
low -= 1
if low < 0:
break
if low <= 0 or index_at(low) != epoch:
raise ServiceError(
503,
"historical epoch boundary is beyond the chain's retained state",
)

while low + 1 < high:
middle = (low + high) // 2
index = self.subtensor.get_subnet_epoch_index(netuid, block=middle)
index = index_at(middle)
if index is None:
raise ServiceError(503, "historical epoch state unavailable")
raise ServiceError(
503,
"historical epoch boundary is beyond the chain's retained state",
)
if index <= epoch:
low = middle
else:
Expand Down Expand Up @@ -411,7 +467,10 @@ async def _loop(self, operation, seconds: float) -> None:
try:
await operation()
except Exception as error:
logging.warning("master background operation failed (%s)", type(error).__name__)
# The type alone is not diagnosable: an epoch loop that fails
# every tick emits nothing, and the only symptom upstream is a
# burn. Keep the message and the traceback.
logging.warning("master background operation failed", exc_info=error)
try:
await asyncio.wait_for(self._stop.wait(), timeout=seconds)
except TimeoutError:
Expand Down
15 changes: 14 additions & 1 deletion src/cortex/protocol/aggregate.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,19 @@ def aggregate_challenge_weights(
kept[key] = score
miner_total = compensated_sum(by_uid.values())
if miner_total <= 1e-12:
# Nothing was claimed. The mass that burns is the allocation no miner
# took, which the challenge document states: the bounty share burns only
# when there is no payable report, and the proof share burns when no
# submission is credited. Burning less than the full proof share here
# would mint emission nobody earned.
# Algorithm 3 scales every share by its claimed score, so an epoch nobody
# scored in declares no mass at all; then the whole vector is the burn.
burn_mass = compensated_sum(fractions.values())
if burn_mass <= 1e-12:
burn_mass = 1.0
# The chain still requires a minimum number of positive weights, so the
# burn is spread over that many uids. It stays a burn either way, but it
# is the declared allocations that decide how much burns.
if max_weight_limit <= 0:
raise ProtocolError(f"max_weight_limit={max_weight_limit} admits no positive weight")
candidates = [0] + sorted(set(hotkey_to_uid.values()) - {0})
Expand All @@ -101,7 +114,7 @@ def aggregate_challenge_weights(
f"max_weight_limit={max_weight_limit}) but only {len(candidates)} "
"usable uid(s) available"
)
by_uid = dict.fromkeys(candidates[:needed], 1.0 / needed)
by_uid = dict.fromkeys(candidates[:needed], burn_mass / needed)
kept = {}
else:
burn = 1.0 - miner_total
Expand Down
58 changes: 48 additions & 10 deletions src/cortex/validator/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@

from .chain import BittensorChain, close_subtensor
from .evidence import peer_app
from .keystore import PrivateKeyWallet, load_private_key_wallet
from .keystore import carries_private_key as _carries_private_key
from .service import SubmissionJournal, TickResult, Validator


Expand Down Expand Up @@ -87,8 +89,10 @@ def parser() -> argparse.ArgumentParser:
arguments.add_argument(
"--consensus-seed-file",
type=Path,
required=True,
help="private sr25519 seed for the same validator hotkey",
help="private sr25519 seed for the same validator hotkey. Required only "
"with --peer-consensus: the seed signs cross-validator root statements "
"and dissents, and the default deployment submits from the gateway "
"without either.",
)
arguments.add_argument(
"--peers", type=Path, help="JSON mapping independent validator hotkeys to HTTPS origins"
Expand Down Expand Up @@ -133,13 +137,23 @@ def load_trust(epoch):
minimum_measurements_version=arguments.minimum_measurements_version,
)

# The seed signs cross-validator root statements and dissents, which only
# exist in a peer-consensus deployment (_crosscheck returns immediately
# otherwise, and _dissent does nothing without a seed). Checking it
# unconditionally refused a plain gateway-backed validator whose hotkey is
# an exported keystore key rather than a mnemonic: such a key carries an
# expanded secret, not the 32-byte seed this check wants, so the check
# failed for a deployment that would never have read the seed at all.
def consensus_seed():
if arguments.consensus_seed_file is None:
raise ProtocolError("--peer-consensus requires --consensus-seed-file")
seed = read_seed(arguments.consensus_seed_file)
if public_key(seed) != wallet.hotkey.public_key:
raise ProtocolError("consensus seed must match validator wallet hotkey")
return seed

consensus_seed()
if arguments.peer_consensus:
consensus_seed()
peers = {}
if arguments.peers:
values = json.loads(arguments.peers.read_text())
Expand Down Expand Up @@ -176,7 +190,10 @@ def consensus_seed():
journal=journal,
http=http,
version_key=arguments.version_key,
consensus_seed=consensus_seed,
# A callable only when peer consensus will read it. Passing one
# unconditionally made `_crosscheck` verify a seed the
# deployment does not use, and refuse every tick over it.
consensus_seed=consensus_seed if arguments.peer_consensus else None,
peers=peers,
peer_consensus=arguments.peer_consensus,
min_peer_sample=arguments.min_peer_sample,
Expand All @@ -198,7 +215,10 @@ def consensus_seed():
result.nonce,
)
except (ProtocolError, httpx.HTTPError, OSError) as error:
logging.warning("validator tick refused (%s)", type(error).__name__)
# The type alone is not diagnosable: a validator that refuses
# every tick submits nothing, and the only upstream symptom is
# a burn. Same reasoning as the master's background loop.
logging.warning("validator tick refused", exc_info=error)
if arguments.once:
raise
if arguments.once:
Expand Down Expand Up @@ -233,11 +253,29 @@ def main(argv: list[str] | None = None) -> None:
subtensor = Subtensor(
network=arguments.network, fallback_endpoints=arguments.fallback_endpoints
)
wallet = Wallet(
name=arguments.wallet_name,
hotkey=arguments.wallet_hotkey,
path=str(Path(arguments.wallet_path).expanduser()),
)
# A hotkey exported from a Polkadot-style keystore has no mnemonic, so the
# wallet library cannot load it (see .hotkey). That file is an alternative
# source, and exactly one source must be given: a wallet whose hotkey is not
# the key the operator means is worse than a refusal at startup.
# A hotkey exported from a Polkadot-style keystore lives at the same path a
# mnemonic wallet does — <path>/<name>/hotkeys/<hotkey> — so the file
# decides how it is read, and the command line stays what the deploy gate
# audits. A mnemonic wallet keeps working exactly as before.
wallet_path = Path(arguments.wallet_path).expanduser()
hotkey_file = wallet_path / arguments.wallet_name / "hotkeys" / arguments.wallet_hotkey
wallet: Wallet | PrivateKeyWallet
try:
private = hotkey_file.is_file() and _carries_private_key(hotkey_file)
if private:
wallet = load_private_key_wallet(hotkey_file)
except ProtocolError as error:
raise SystemExit(f"validator wallet refused: {error}") from None
if not private:
wallet = Wallet(
name=arguments.wallet_name,
hotkey=arguments.wallet_hotkey,
path=str(wallet_path),
)
try:
result = asyncio.run(run(arguments, subtensor, wallet))
if (
Expand Down
Loading
Loading