From a4f3ea4ac12dd720bbaa1567ba078223441b66a6 Mon Sep 17 00:00:00 2001 From: echobt <154886644+echobt@users.noreply.github.com> Date: Mon, 28 Sep 2026 20:50:09 +0000 Subject: [PATCH] feat(master): send the completed epoch's chain time to challenges The master reads Timestamp.Now at the pinned end-block hash of the completed epoch (integer milliseconds, floored to seconds) and passes it to challenge containers as the authenticated get_weights query epoch_at. OpenType prices its champion decay by this chain time, not request time. - ChainSnapshot.timestamp_seconds is optional; an unavailable timestamp leaves validator snapshots unchanged (reorg check still applies). - Under algorithm 3 with trusted OpenType, a missing or invalid timestamp postpones emission before any leaf is replaced or the epoch is sealed, instead of signing NoScore. Bounty/proof-only setups are unaffected. No wall-clock or estimated block-time fallback. - No signed protocol or validator economics change. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/CHALLENGES.md | 10 ++ src/cortex/challenges/client.py | 13 ++- src/cortex/master.py | 17 +++- src/cortex/validator/chain.py | 13 +++ src/cortex/validator/service.py | 1 + tests/challenges/test_contract.py | 38 +++++++- tests/test_master.py | 154 +++++++++++++++++++++++++++++- tests/validator/test_chain.py | 56 +++++++++++ 8 files changed, 295 insertions(+), 7 deletions(-) diff --git a/docs/CHALLENGES.md b/docs/CHALLENGES.md index 0268faa67..6ff09099d 100644 --- a/docs/CHALLENGES.md +++ b/docs/CHALLENGES.md @@ -55,6 +55,16 @@ receives a leaf-signing seed. `capabilities` is a subset of `get_weights` and `proxy_routes`. +The master also sends `epoch_at=` when available: +`Timestamp.Now` at the completed epoch's pinned end-block hash, validated as +integer milliseconds and floored to seconds. There is no wall-clock or estimated +block-time fallback. Algorithm 3 epochs with trusted OpenType require this time; +missing or invalid time postpones emission before any leaves are replaced or the +epoch is sealed, rather than signing `NoScore`. Other challenge configurations +and legacy client calls remain compatible without `epoch_at`. Containers that do +not use epoch time may ignore this optional query parameter. Persisted epoch +responses remain immutable regardless of subsequent request parameters. + `get_weights` requires `Authorization: Bearer ` (`401` otherwise) and `X-Platform-Challenge-Slug: ` (`403` on mismatch). It returns: diff --git a/src/cortex/challenges/client.py b/src/cortex/challenges/client.py index 449bf659d..1a44cad06 100644 --- a/src/cortex/challenges/client.py +++ b/src/cortex/challenges/client.py @@ -15,6 +15,7 @@ from cortex.http import read_private_file from cortex.protocol.crypto import decode_hotkey from cortex.protocol.models import FULL_SHARE_SCORE, NoScore, NoScoreReason, Score +from cortex.protocol.scale import uint from .registry import RegistryEntry @@ -88,14 +89,22 @@ class ChallengeClient: def __init__(self, http: httpx.AsyncClient, secrets_dir: Path, *, retry_seconds: float = 5): self.http, self.secrets_dir, self.retry_seconds = http, secrets_dir, retry_seconds - async def weights(self, entry: RegistryEntry, epoch: int) -> ChallengeWeights: + async def weights( + self, entry: RegistryEntry, epoch: int, *, epoch_at: int | None = None + ) -> ChallengeWeights: + params = {"epoch": str(epoch)} + if epoch_at is not None: + uint(epoch_at, 8) + if epoch_at == 0: + raise ValueError("epoch_at must be positive") + params["epoch_at"] = str(epoch_at) token = read_private_file(self.secrets_dir / entry.id / "internal.token") for attempt in range(ATTEMPTS): try: async with self.http.stream( "GET", f"{entry.url}/internal/v1/get_weights", - params={"epoch": str(epoch)}, + params=params, headers={ "authorization": f"Bearer {token}", "x-platform-challenge-slug": entry.id, diff --git a/src/cortex/master.py b/src/cortex/master.py index f2ef2284a..dd09c9121 100644 --- a/src/cortex/master.py +++ b/src/cortex/master.py @@ -333,7 +333,9 @@ def __init__( self.challenge_seed, self.registry, self.challenges = challenge_seed, registry, challenges self._lock = asyncio.Lock() - async def _scores(self, challenge: bytes, epoch: int, expected: set[bytes]): + async def _scores( + self, challenge: bytes, epoch: int, expected: set[bytes], epoch_at: int | None = None + ): algorithm = self.gateway.trust.algorithm_version try: if challenge != b"proof": @@ -342,7 +344,7 @@ async def _scores(self, challenge: bytes, epoch: int, expected: set[bytes]): raise ServiceError(503, "challenge container not registered") if algorithm == 1: raise ServiceError(503, "container challenges need algorithm 2 or 3") - answer = await self.challenges.weights(entry, epoch) + answer = await self.challenges.weights(entry, epoch, epoch_at=epoch_at) return leaf_scores(answer, expected, algorithm_version=algorithm) # Backend readiness is still required before old rows can be emitted. readiness = await self.proof.backend.readiness() @@ -374,6 +376,15 @@ async def _scores(self, challenge: bytes, epoch: int, expected: set[bytes]): async def _emit(self, epoch: int, snapshot: ChainSnapshot) -> None: self.gateway.refresh_trust(epoch) + epoch_at = snapshot.timestamp_seconds + if type(epoch_at) is not int or not 0 < epoch_at < 2**64: + epoch_at = None + if ( + self.gateway.trust.algorithm_version == 3 + and any(entry.id == b"opentype" for entry in self.gateway.trust.challenges) + and epoch_at is None + ): + raise ServiceError(503, "completed epoch timestamp unavailable") self.proof.topic_public_key = next( (entry.public_key for entry in self.gateway.trust.challenges if entry.id == b"proof"), self.proof.topic_public_key, @@ -384,7 +395,7 @@ async def _emit(self, epoch: int, snapshot: ChainSnapshot) -> None: seed = self.challenge_seed(challenge.id) if public_key(seed) != challenge.public_key: raise ServiceError(503, "challenge signing key does not match owner trust") - outcomes = await self._scores(challenge.id, epoch, expected) + outcomes = await self._scores(challenge.id, epoch, expected, epoch_at) # No unknown score key can expand E; missing scores are explicit signed absences. leaves = [] for hotkey in sorted(expected): diff --git a/src/cortex/validator/chain.py b/src/cortex/validator/chain.py index 607783a5c..f1554bf33 100644 --- a/src/cortex/validator/chain.py +++ b/src/cortex/validator/chain.py @@ -102,6 +102,18 @@ def _snapshot(self, block: int, netuid: int) -> ChainSnapshot: owner = self.subtensor.get_subnet_owner_hotkey(netuid, block=block) epoch = self.subtensor.get_subnet_epoch_index(netuid, block=block) uint(epoch, 8) + timestamp_seconds = None + try: + timestamp = self.subtensor.substrate.query( + module="Timestamp", storage_function="Now", params=[], block_hash=before + ).value + uint(timestamp, 8) + if timestamp < 1000: + raise ProtocolError("invalid chain timestamp") + timestamp_seconds = timestamp // 1000 + except Exception: + # Validators do not need time; the master gates time-dependent emission. + logging.warning("historical chain timestamp unavailable block=%d", block) after = self.subtensor.get_block_hash(block) if before != after: raise ProtocolError("chain reorganized during metagraph read") @@ -126,6 +138,7 @@ def _snapshot(self, block: int, netuid: int) -> ChainSnapshot: decode_hotkey(owner), frozenset(permits), epoch, + timestamp_seconds, ) async def submit( diff --git a/src/cortex/validator/service.py b/src/cortex/validator/service.py index 43ed08df1..5c439d873 100644 --- a/src/cortex/validator/service.py +++ b/src/cortex/validator/service.py @@ -53,6 +53,7 @@ class ChainSnapshot: owner_hotkey: bytes validator_permits: frozenset[int] epoch: int | None = None + timestamp_seconds: int | None = None class Chain(Protocol): diff --git a/tests/challenges/test_contract.py b/tests/challenges/test_contract.py index f18b935ff..dc43c21cf 100644 --- a/tests/challenges/test_contract.py +++ b/tests/challenges/test_contract.py @@ -7,7 +7,7 @@ from fastapi import FastAPI from fastapi.testclient import TestClient -from cortex.challenges.client import ChallengeWeights, leaf_scores, parse_weights +from cortex.challenges.client import ChallengeClient, ChallengeWeights, leaf_scores, parse_weights from cortex.challenges.proxy import create_router from cortex.challenges.registry import parse_registry from cortex.protocol.models import FULL_SHARE_SCORE, NoScore, NoScoreReason, Score @@ -92,6 +92,42 @@ def test_weights_for_another_challenge_epoch_or_with_invalid_values_are_rejected parse_weights(body, slug="opentype", epoch=7) +@pytest.mark.parametrize("epoch_at", [None, 1_790_000_123]) +async def test_weights_transports_optional_epoch_at_with_authentication(tmp_path, epoch_at): + entry = parse_registry({"version": 1, "challenge": [ROW]})["opentype"] + secret = tmp_path / "opentype" / "internal.token" + secret.parent.mkdir() + secret.write_text("fixture-internal") + secret.chmod(0o600) + seen = [] + + def upstream(request): + seen.append(request) + return httpx.Response(200, json={"challenge_slug": "opentype", "epoch": 7, "weights": {}}) + + async with httpx.AsyncClient(transport=httpx.MockTransport(upstream)) as http: + answer = await ChallengeClient(http, tmp_path).weights(entry, 7, epoch_at=epoch_at) + assert answer.weights == {} + assert len(seen) == 1 + assert dict(seen[0].url.params) == ( + {"epoch": "7"} if epoch_at is None else {"epoch": "7", "epoch_at": str(epoch_at)} + ) + assert seen[0].headers["authorization"] == "Bearer fixture-internal" + assert seen[0].headers["x-platform-challenge-slug"] == "opentype" + + +@pytest.mark.parametrize("epoch_at", [True, False, 0, -1, 1.5, "1790000123", 2**64]) +async def test_weights_rejects_invalid_epoch_at_before_http_or_secrets(tmp_path, epoch_at): + entry = parse_registry({"version": 1, "challenge": [ROW]})["opentype"] + + def upstream(request): + pytest.fail("invalid epoch_at must never reach HTTP") + + async with httpx.AsyncClient(transport=httpx.MockTransport(upstream)) as http: + with pytest.raises(ValueError): + await ChallengeClient(http, tmp_path).weights(entry, 7, epoch_at=epoch_at) + + def test_proxy_forwards_public_routes_only(): seen = [] diff --git a/tests/test_master.py b/tests/test_master.py index 2dde1217b..0dc475913 100644 --- a/tests/test_master.py +++ b/tests/test_master.py @@ -28,6 +28,7 @@ def __init__(self): ) self.submissions = [] self.fail = False + self.timestamp_seconds = 1_790_000_123 async def epoch_state(self, netuid): if self.fail: @@ -39,7 +40,12 @@ async def current_block(self): async def snapshot(self, block, netuid): return ChainSnapshot( - block, sha256(str(block).encode()).digest(), self.rows, self.rows[0].hotkey, frozenset() + block, + sha256(str(block).encode()).digest(), + self.rows, + self.rows[0].hotkey, + frozenset(), + timestamp_seconds=self.timestamp_seconds, ) async def submit(self, netuid, vector, version_key): @@ -442,3 +448,149 @@ def test_backend_config_rejects_credentials_in_url(tmp_path): def test_config_rejects_invalid_emission_intervals(tmp_path, value): with pytest.raises(ValueError, match="intervals"): replace(master_config(tmp_path), emit_poll_seconds=value) + + +@pytest.fixture +async def decay_master(tmp_path): + from cortex.challenges.registry import parse_registry + + config = master_config(tmp_path) + key = tmp_path / "opentype.key" + key.write_bytes(bytes([3]) * 32) + key.chmod(0o600) + entries = [] + for slug in ("bounty", "opentype"): + token = config.challenge_secrets_dir / slug / "internal.token" + token.parent.mkdir(parents=True) + token.write_text(f"{slug}-internal") + token.chmod(0o600) + entries.append( + { + "id": slug, + "image": f"ghcr.io/fixture/{slug}", + "source": f"https://github.com/fixture/{slug}", + } + ) + chain = FakeEpochChain() + trust = TrustRoot( + ( + ChallengeEntry(b"bounty", public_key(bytes([1]) * 32), 3000), + ChallengeEntry(b"opentype", public_key(bytes([3]) * 32), 7000), + ), + sha256(b"\x00").digest(), + public_key(bytes([7]) * 32), + challenges_version=3, + ) + seen = [] + + def upstream(request): + slug = request.headers["x-platform-challenge-slug"] + assert request.headers["authorization"] == f"Bearer {slug}-internal" + seen.append(request) + return httpx.Response( + 200, + json={ + "challenge_slug": slug, + "epoch": 12, + "weights": {chain.rows[1].hotkey.hex(): 0.5} if slug == "opentype" else {}, + "full_share_mass": 1.0, + }, + ) + + runtime = await build_master( + config, + chain=chain, + epochs=chain, + trust=trust, + challenge_http=httpx.AsyncClient(transport=httpx.MockTransport(upstream)), + ) + runtime.emitter.registry = lambda: parse_registry({"version": 1, "challenge": entries}) + try: + yield SimpleNamespace(runtime=runtime, chain=chain, trust=trust, seen=seen) + finally: + await runtime.close() + + +@pytest.mark.parametrize("timestamp", [None, True, False, 0, -1, 1.5, "1790000123", 2**64]) +async def test_missing_or_invalid_epoch_time_postpones_before_any_leaves_or_seal( + decay_master, timestamp +): + from cortex.protocol import sign_leaf + + master = decay_master + runtime = master.runtime + await runtime.emitter.tick() + original = sign_leaf(bytes([1]) * 32, b"bounty", master.chain.rows[1].hotkey, 12, Score(123)) + runtime.gateway.accept_leaf(original) + before = runtime.gateway.store.leaves(12) + master.chain.timestamp_seconds = timestamp + master.chain.state = EpochState(13, 100, 105) + with pytest.raises(ServiceError, match="epoch timestamp unavailable"): + await runtime.emitter.tick() + assert runtime.gateway.store.leaves(12) == before + assert runtime.gateway.store.bundle(12) is None + assert master.seen == [] + assert runtime.emitter.journal.pending(13) + master.chain.timestamp_seconds = 1_790_000_123 + assert await runtime.emitter.tick() == [12] + assert runtime.emitter.journal.pending(13) == [] + + +async def test_historical_epoch_time_reaches_http_and_decayed_mass_seals_without_renormalization( + decay_master, +): + master = decay_master + runtime = master.runtime + await runtime.emitter.tick() + master.chain.state = EpochState(13, 100, 105) + assert await runtime.emitter.tick() == [12] + assert all( + dict(request.url.params) == {"epoch": "12", "epoch_at": "1790000123"} + for request in master.seen + ) + assert len(master.seen) == 2 + leaves = runtime.gateway.store.leaves(12) + assert next( + leaf.score + for leaf in leaves + if leaf.challenge_id == b"opentype" and leaf.miner_hotkey == master.chain.rows[1].hotkey + ) == Score(500_000_000_000) + latest = runtime.gateway.latest() + assert latest["sealed"] is True + assert latest["final_vector"] == [[0, 42598], [1, 22937]] + original = runtime.gateway.bundle_bytes(12) + master.chain.timestamp_seconds = None + async with httpx.AsyncClient( + transport=httpx.ASGITransport(runtime.app()), base_url="https://master" + ) as http: + journal = SubmissionJournal(":memory:") + try: + validator = Validator( + gateway_url="https://master", + netuid=541, + trust=master.trust, + chain=master.chain, + journal=journal, + http=http, + ) + assert (await validator.run_once()).outcome == "submitted" + finally: + journal.close() + assert master.chain.submissions == [((0, 42598), (1, 22937))] + assert await runtime.emitter.tick() == [] + assert runtime.gateway.bundle_bytes(12) == original + + +async def test_bounty_only_algorithm_three_still_seals_without_epoch_time(decay_master): + master = decay_master + runtime = master.runtime + runtime.gateway.trust = replace( + master.trust, challenges=(ChallengeEntry(b"bounty", public_key(bytes([1]) * 32), 10000),) + ) + master.chain.timestamp_seconds = None + await runtime.emitter.tick() + master.chain.state = EpochState(13, 100, 105) + assert await runtime.emitter.tick() == [12] + assert len(master.seen) == 1 + assert dict(master.seen[0].url.params) == {"epoch": "12"} + assert runtime.gateway.latest()["sealed"] is True diff --git a/tests/validator/test_chain.py b/tests/validator/test_chain.py index a12df0b10..a4ef354a2 100644 --- a/tests/validator/test_chain.py +++ b/tests/validator/test_chain.py @@ -457,3 +457,59 @@ def epoch_index(netuid, *, block): ("get_subnet_owner_hotkey", 541, 99), ("get_subnet_epoch_index", 541, 99), ] + + +@pytest.mark.parametrize( + "timestamp,seconds", + [ + (1_790_000_123_999, 1_790_000_123), + (1000, 1), + (2**64 - 1, (2**64 - 1) // 1000), + (None, None), + (True, None), + (0, None), + (999, None), + (-1, None), + (1.5, None), + ("1790000123000", None), + (2**64, None), + ], +) +async def test_historical_timestamp_is_validated_and_read_at_pinned_hash(timestamp, seconds): + digest = "0x" + "04" * 32 + + def query(*, module, storage_function, params, block_hash): + assert (module, storage_function, params, block_hash) == ("Timestamp", "Now", [], digest) + return SimpleNamespace(value=timestamp) + + subtensor = SimpleNamespace( + get_block_hash=lambda block: digest, + neurons_lite=lambda netuid, *, block: [], + get_subnet_owner_hotkey=lambda netuid, *, block: encode_hotkey(bytes([1]) * 32), + get_subnet_epoch_index=lambda netuid, *, block: 12, + substrate=SimpleNamespace(query=query), + ) + snapshot = await BittensorChain(subtensor, None).snapshot(99, 541) + assert snapshot.timestamp_seconds == seconds + assert snapshot.block_hash == bytes.fromhex(digest[2:]) + + +@pytest.mark.parametrize("reorg", [False, True]) +async def test_timestamp_rpc_failure_preserves_validator_snapshot_but_not_reorg(reorg): + hashes = iter(["0x" + "01" * 32, "0x" + ("02" if reorg else "01") * 32]) + + def query(**kwargs): + raise OSError("historical state unavailable") + + subtensor = SimpleNamespace( + get_block_hash=lambda block: next(hashes), + neurons_lite=lambda netuid, *, block: [], + get_subnet_owner_hotkey=lambda netuid, *, block: encode_hotkey(bytes([1]) * 32), + get_subnet_epoch_index=lambda netuid, *, block: 12, + substrate=SimpleNamespace(query=query), + ) + if reorg: + with pytest.raises(ProtocolError, match="reorganized"): + await BittensorChain(subtensor, None).snapshot(99, 541) + else: + assert (await BittensorChain(subtensor, None).snapshot(99, 541)).timestamp_seconds is None