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
10 changes: 10 additions & 0 deletions docs/CHALLENGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -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=<positive integer Unix seconds>` 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 <internal token>` (`401`
otherwise) and `X-Platform-Challenge-Slug: <slug>` (`403` on mismatch). It
returns:
Expand Down
13 changes: 11 additions & 2 deletions src/cortex/challenges/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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,
Expand Down
17 changes: 14 additions & 3 deletions src/cortex/master.py
Original file line number Diff line number Diff line change
Expand Up @@ -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":
Expand All @@ -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()
Expand Down Expand Up @@ -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,
Expand All @@ -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):
Expand Down
13 changes: 13 additions & 0 deletions src/cortex/validator/chain.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -126,6 +138,7 @@ def _snapshot(self, block: int, netuid: int) -> ChainSnapshot:
decode_hotkey(owner),
frozenset(permits),
epoch,
timestamp_seconds,
)

async def submit(
Expand Down
1 change: 1 addition & 0 deletions src/cortex/validator/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
38 changes: 37 additions & 1 deletion tests/challenges/test_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 = []

Expand Down
154 changes: 153 additions & 1 deletion tests/test_master.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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):
Expand Down Expand Up @@ -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
Loading
Loading