diff --git a/Documentation/public/roles/auditors.md b/Documentation/public/roles/auditors.md index b53fbc9..8703d1d 100644 --- a/Documentation/public/roles/auditors.md +++ b/Documentation/public/roles/auditors.md @@ -89,7 +89,7 @@ dincli auditor lms-evaluation evaluate [--gi ] [--submit] [ | Option | Required | Description | |---|---|---| | `--gi ` | No | Global Iteration index. Defaults to the current GI if omitted | -| `--submit` | No | Submit evaluation scores and eligibility flags to the blockchain | +| `--submit` | No | **Commit** a hidden hash of each score and eligibility vote on-chain (the score, vote and a random salt are cached locally for the reveal). Omit to evaluate locally only. Re-running skips models you've already committed | | `--batch ` | No | Evaluate a specific batch. Defaults to your assigned batch if omitted | | `--lmi ` | No | Evaluate a single Local Model by index. Evaluates all models in the batch if omitted | diff --git a/Documentation/technical/contracts/DINTaskAuditor.md b/Documentation/technical/contracts/DINTaskAuditor.md index f53398c..8fe0972 100644 --- a/Documentation/technical/contracts/DINTaskAuditor.md +++ b/Documentation/technical/contracts/DINTaskAuditor.md @@ -203,7 +203,7 @@ Read alongside the [foundry/src security review](../audits/foundry-src-security- - **No. 6 — Stale NatSpec and reused errors:** `slashAuditors` says S1 and S3 are both `minStake()` (S1 is now a fraction). Several comments call parameters "DAO-settable"; they are `onlyOwner`, i.e. set by the model owner. `setDisputePenaltyBps` reuses `TA_InvalidDisputeBond`, and `closeExpiredDispute` reuses `TA_DisputeWindowClosed` for "window still open". - **No. 7 — `modelId` is fixed at construction**, before the registry assigns it (see [DINTaskCoordinator §10 No. 4](DINTaskCoordinator.md#10-review-notes--open-caveats)). - **No. 8 — dincli lags this contract:** `dincli model-owner deploy task-auditor` still calls the older two-argument constructor (no `modelId`). -- **No. 9 — dincli auditor commit retry can lose the committed salt.** Rerunning `dincli auditor lms-evaluation evaluate --submit` generates a new salt and overwrites the local commit cache even when the on-chain commit already exists. The later reveal then fails the hash check and the auditor is slashed for a missed vote. Tracked in issue No. 202 (Part 1); until it is fixed, don't rerun the command for a GI that already has commits. +- **No. 9 — Fixed: a dincli auditor commit retry no longer loses the committed salt.** Rerunning `dincli auditor lms-evaluation evaluate --submit` used to generate a new salt and overwrite the local commit cache even when the on-chain commit already existed. The later reveal then failed the hash check, and the auditor was slashed for a missed vote. dincli now skips any LM whose `hasCommittedLM` is already set, leaving its cache untouched, and writes the cache before sending the commit tx (issue No. 202 Part 1). A rerun is safe. --- diff --git a/Documentation/technical/contracts/DINTaskCoordinator.md b/Documentation/technical/contracts/DINTaskCoordinator.md index 141e9c1..1b05b7c 100644 --- a/Documentation/technical/contracts/DINTaskCoordinator.md +++ b/Documentation/technical/contracts/DINTaskCoordinator.md @@ -279,7 +279,7 @@ Earlier findings from the [foundry/src security review](../audits/foundry-src-se - **No. 7 — Stale NatSpec:** comments around the dispute scaffold still say `DinTreasury` "doesn't exist on develop yet" (forfeitures are already forwarded). Several parameters are described as "DAO-settable"; they are `onlyOwner`, i.e. set by the model owner. - **No. 8 — Leftovers:** `networkFeeFloor` is stored but not enforced. `setTestDataAssignedFlag` gates nothing: evaluation can start without test data being assigned. `releaseGIRegistrationSlots` uses string `require` messages, unlike the rest of the contract. - **No. 9 — Not upgradeable:** a bug in a model's task contracts requires redeploying them and re-registering the model. -- **No. 10 — dincli lags this contract:** `dincli model-owner deploy task-coordinator` still calls the older one-argument constructor (no `modelId`). `dincli aggregator aggregate-t2` names its working directory, worker job and container after the last T1 batch id, not the T2 batch id (issue #202, Part 2); the on-chain commit is unaffected. +- **No. 10 — dincli lags this contract:** `dincli model-owner deploy task-coordinator` still calls the older one-argument constructor (no `modelId`). (`dincli aggregator aggregate-t2` used to name its working directory, worker job and container after the last T1 batch id; fixed in issue No. 202 Part 2.) - **No. 11 — Fixed: `registerDINaggregator` now checks the GI.** It used to have no `onlyCurrentGI`, so during GI N's window a validator could register for GI N+1. Up to 300 addresses could fill that list, which locked out honest registrants and handed the attacker every T1/T2 batch. Registering for a released past GI also leaked the caller's own slot. Fixed by adding `onlyCurrentGI` (issue #206); it cost 9 bytes inside the No. 1 budget. --- diff --git a/dincli/cli/aggregator.py b/dincli/cli/aggregator.py index 8b76214..e41851f 100644 --- a/dincli/cli/aggregator.py +++ b/dincli/cli/aggregator.py @@ -513,12 +513,12 @@ def aggregate_t2( found_batch = False - if batch_id and batch_id >= t2_batches_count: + if batch_id is not None and batch_id >= t2_batches_count: console.print(f"[red]Error:[/red] invalid T2 batch ID {batch_id} does not exist") raise typer.Exit(1) for i in range(t2_batches_count): - if batch_id: + if batch_id is not None: if i != batch_id: continue @@ -544,8 +544,11 @@ def aggregate_t2( for j in range(t1_batches_count): time.sleep(0.1) - (bid, val, idxs, fin, cid) = taskCoordinator_contract.functions.getTier1Batch(curr_GI, j).call() - model_cids.append(get_cid_from_bytes32(cid.hex())) + # Distinct names: rebinding `bid` here used to leave the last T1 + # batch id in it, mis-naming the T2 models path, worker job and + # container below (issue #202 Part 2). + (_, _, _, _, t1_final_cid) = taskCoordinator_contract.functions.getTier1Batch(curr_GI, j).call() + model_cids.append(get_cid_from_bytes32(t1_final_cid.hex())) console.print(f"Aggregating T2 batch {bid} for aggregator {account.address} with T1 final cids {model_cids} and genesis model cid {genesis_model_ipfs_hash}") diff --git a/dincli/cli/auditor.py b/dincli/cli/auditor.py index 7e30391..01697b8 100644 --- a/dincli/cli/auditor.py +++ b/dincli/cli/auditor.py @@ -430,6 +430,15 @@ def evaluate_lms( continue found_any = True + + # A retry must never replace the salt behind an existing on-chain + # commitment -- the reveal would then fail TA_RevealHashMismatch + # and the auditor be S1-slashed for an honest vote (issue #202). + # Skip before re-evaluating, leaving the commit cache untouched. + if submit and task_auditor_contract.functions.hasCommittedLM(curr_GI, batch_id, account.address, model_index).call(): + console.print(f"[yellow]LM {model_index} in audit batch {batch_id} already committed by {account.address}; skipping (reveal with `dincli auditor lms-evaluation reveal`).[/yellow]") + continue + console.print(f"[bold green]Evaluating LM {model_index} from Audit batch {batch_id}![/bold green]") time.sleep(0.5) @@ -550,6 +559,13 @@ def evaluate_lms( try: time.sleep(0.5) + # Cache before sending: build_and_send_tx returns None both + # on a revert and when the receipt wait fails for a tx that + # may still mine, so saving only after the send could strand + # a real commit without its preimage. A failed attempt leaves + # a stale cache that the next retry replaces (the LM isn't + # committed yet, so the hasCommittedLM guard lets it through). + _save_commit(model_base_dir, curr_GI, batch_id, model_index, score_int, vote_bool, salt) build_and_send_tx( ctx, task_auditor_contract.functions.commitAuditScore(curr_GI, batch_id, model_index, commit_hash), @@ -558,9 +574,6 @@ def evaluate_lms( f"Audit score commit failed for LM {model_index} from batch {batch_id}!", exit_on_failure=False ) - # Only cache locally once the commit tx is known to have - # been attempted -- reveal is a no-op without this file. - _save_commit(model_base_dir, curr_GI, batch_id, model_index, score_int, vote_bool, salt) except Exception as e: console.print(f"[bold red]✗ Error committing audit score for LM {model_index} from batch {batch_id}: {e}[/bold red]") diff --git a/tests/test_aggregator_t2_batch_id.py b/tests/test_aggregator_t2_batch_id.py new file mode 100644 index 0000000..88f63d9 --- /dev/null +++ b/tests/test_aggregator_t2_batch_id.py @@ -0,0 +1,105 @@ +"""`aggregator aggregate-t2` names its work after the T2 batch (issue #202 Part 2). + +aggregate_t2 reads the T2 batch id from getTier2Batch, then loops over every +T1 batch to collect their final CIDs. That loop used to unpack into `bid` +too, so after it `bid` held the *last T1 batch's* id, and the T2 models path, +worker job and container were named after an unrelated T1 batch. The on-chain +commit used the loop index and was unaffected. +""" +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pytest + +from dincli.cli import aggregator +from dincli.cli.aggregator import aggregate_t2 + +GI = 1 +ACCOUNT = "0x" + "11" * 20 +T1_BATCH_COUNT = 3 # T1 ids 0..2; the last one (2) is what leaked before the fix + + +class _Call: + def __init__(self, value): + self._value = value + + def call(self): + return self._value + + +def _make_contract(): + def getter(name): + if name == "genesisModelIpfsHash": + return lambda: _Call(b"\x00" * 32) + if name == "getTier2Batch": + return lambda gi, i: _Call((0, [ACCOUNT], False, b"\x00" * 32)) + if name == "getAggregatorSubmission": + return lambda gi, tier, bid, who: _Call((False, b"\x00" * 32, False, b"\x00" * 32, 0)) + if name == "tier1BatchCount": + return lambda gi: _Call(T1_BATCH_COUNT) + if name == "getTier1Batch": + return lambda gi, j: _Call((j, ["0x" + "44" * 20], [j], True, bytes([j + 1]) * 32)) + raise AttributeError(name) + + class _Functions: + def __getattr__(self, name): + return getter(name) + + contract = MagicMock() + contract.functions = _Functions() + return contract + + +@pytest.fixture +def ctx(tmp_path): + ctx = MagicMock() + account = MagicMock() + account.address = ACCOUNT + ctx.obj.get_en_w3_account_console.return_value = ("local", MagicMock(), account, MagicMock()) + ctx.obj.get_deployed_din_task_coordinator_contract.return_value = _make_contract() + ctx.obj.get_current_gi_and_state.return_value = (GI, 0) + ctx.obj.get_model_base_dir.return_value = tmp_path + return ctx + + +@pytest.fixture +def stubs(tmp_path): + with patch.object(aggregator, "get_cid_from_bytes32", side_effect=lambda h: f"Qm{h[-4:]}"), \ + patch.object(aggregator, "get_manifest_key", + return_value={"path": "services/aggregator.py", "ipfs": "QmSvc"}), \ + patch.object(aggregator, "require_custom_manifest_service"), \ + patch.object(aggregator, "write_worker_job", + return_value=(tmp_path / "job.json", tmp_path / "out")) as write_job, \ + patch.object(aggregator, "run_worker_container", + return_value=MagicMock(returncode=0, stdout="", stderr="")) as run_worker, \ + patch.object(aggregator, "read_worker_result", + return_value={"status": "ok", "result": "/din/model/avg.pth"}), \ + patch.object(aggregator, "upload_to_ipfs", return_value="QmAgg"), \ + patch.object(aggregator.time, "sleep"): + yield write_job, run_worker + + +@pytest.mark.parametrize("batch_id", [None, 0]) +def test_t2_work_is_named_after_the_t2_batch_not_the_last_t1_batch(ctx, stubs, tmp_path, batch_id): + write_job, run_worker = stubs + + # aggregate_t2(ctx, model_id, gi, submit, batch_id, packages_dir, no_cache) + aggregate_t2(ctx, 1, None, False, batch_id, None, False) + + write_job.assert_called_once() + assert write_job.call_args.args[1] == f"aggregator_t2_gi_{GI}_batch_0" + job_args = write_job.call_args.args[2]["args"] + assert job_args[4] == 0 # bid passed to get_aggregated_cid_t2 + assert len(job_args[2]) == T1_BATCH_COUNT # every T1 final CID collected + + run_worker.assert_called_once() + assert run_worker.call_args.kwargs["container_name"].endswith(f"-gi-{GI}-batch-0") + models_path = Path(run_worker.call_args.kwargs["writable_subdirs"][0]) + assert models_path == tmp_path / "aggregator" / ACCOUNT / str(GI) / "T2" / "0" / "models" + + +def test_batch_out_of_range_still_rejected(ctx, stubs): + import typer + + with pytest.raises(typer.Exit): + aggregate_t2(ctx, 1, None, False, 1, None, False) diff --git a/tests/test_auditor_commit_retry.py b/tests/test_auditor_commit_retry.py new file mode 100644 index 0000000..c2cd4aa --- /dev/null +++ b/tests/test_auditor_commit_retry.py @@ -0,0 +1,187 @@ +"""Retry safety for `auditor lms-evaluation evaluate --submit` (issue #202 Part 1). + +The auditor-side twin of the aggregator retry guard (PR No. 197 review, +finding No. 1; see tests/test_aggregator_commit_retry.py). The commit-then- +reveal flow caches the committed (score, vote, salt) per LM so the later +reveal can reproduce the commit hash. Two properties keep a retry from +stranding an honest auditor, whose reveal would otherwise revert +TA_RevealHashMismatch and get them S1-slashed (AUD_NO_VOTE): + +- an LM already committed on-chain (hasCommittedLM) is skipped before + re-evaluating, and its cached (score, vote, salt) is left untouched; +- for an uncommitted LM, the cache is written *before* the commit tx is sent, + so a tx whose receipt wait fails but which still mines never loses its + preimage (build_and_send_tx returns None in that case, same as on a revert). +""" +import json +from unittest.mock import MagicMock, patch + +import pytest + +from dincli.cli import auditor +from dincli.cli.auditor import (_audit_commit_hash, _commit_store_path, + _save_commit, evaluate_lms) + +GI = 1 +BATCH = 0 +ACCOUNT = "0x" + "11" * 20 +OWNER = "0x" + "22" * 20 +MODELS = [0, 1] +OLD_SALT = b"\xee" * 32 + + +class _Call: + def __init__(self, value): + self._value = value + + def call(self): + return self._value + + +def _make_auditor_contract(state): + """state: 'committed' (set of model indexes), 'commits' (list of args).""" + + def getter(name): + if name == "AuditorsBatchCount": + return lambda gi: _Call(1) + if name == "getAuditorsBatch": + return lambda gi, b: _Call((BATCH, [ACCOUNT], MODELS, b"\x01" * 16)) + if name == "hasCommittedLM": + return lambda gi, b, who, m: _Call(m in state["committed"]) + if name == "lmSubmissions": + return lambda gi, m: _Call(("0x" + "33" * 20, b"\x00" * 32, 0, False, False, False, 0)) + if name == "encryptedTestDataKey": + return lambda gi, b, who: _Call(b"\x02" * 16) + if name == "commitAuditScore": + def _commit(*args): + state["commits"].append(args) + return MagicMock(name="commitAuditScore_tx") + return _commit + raise AttributeError(name) + + class _Functions: + def __getattr__(self, name): + return getter(name) + + contract = MagicMock() + contract.functions = _Functions() + return contract + + +def _make_coordinator_contract(): + contract = MagicMock() + contract.functions.genesisModelIpfsHash.return_value = _Call(b"\x00" * 32) + contract.functions.owner.return_value = _Call(OWNER) + return contract + + +def _make_ctx(auditor_contract, model_base_dir): + ctx = MagicMock() + account = MagicMock() + account.address = ACCOUNT + ctx.obj.get_en_w3_account_console.return_value = ("local", MagicMock(), account, MagicMock()) + ctx.obj.get_deployed_din_task_coordinator_contract.return_value = _make_coordinator_contract() + ctx.obj.get_deployed_din_task_auditor_contract.return_value = auditor_contract + ctx.obj.get_current_gi_and_state.return_value = (GI, 0) + ctx.obj.get_model_base_dir.return_value = model_base_dir + return ctx + + +def _manifest_key(network, key, model_id): + if key == "requirements.txt": + return {} + if key == "owner_encryption_pubkey": + return "00" * 32 + return {"path": f"services/{key}.py", "ipfs": "QmSvc"} + + +@pytest.fixture +def worker_stubs(tmp_path): + """Stubs the manifest/decrypt/docker steps between LM selection and the + commit, so the evaluation path runs through to the commit tx with + score 80 / eligible True.""" + pt_path = tmp_path / "dataset" / "auditor" / "TestDatasets" / f"auditorDataset_{GI}_{BATCH}.pt" + pt_path.parent.mkdir(parents=True) + pt_path.write_bytes(b"test data") # already decrypted: skips the IPFS fetch + with patch.object(auditor, "get_cid_from_bytes32", return_value="QmCid"), \ + patch.object(auditor, "get_manifest_key", side_effect=_manifest_key), \ + patch.object(auditor, "require_custom_manifest_service"), \ + patch.object(auditor, "_load_auditor_x25519_key"), \ + patch.object(auditor, "PublicKey"), \ + patch.object(auditor, "Box"), \ + patch.object(auditor, "_decrypt_aes_gcm", return_value=b"\x03" * 32 + b"\x04" * 65), \ + patch.object(auditor.Account, "recover_message", return_value=OWNER), \ + patch.object(auditor, "write_worker_job", + return_value=(tmp_path / "job.json", tmp_path / "out")), \ + patch.object(auditor, "run_worker_container", + return_value=MagicMock(returncode=0, stdout="", stderr="")) as run_worker, \ + patch.object(auditor, "read_worker_result", + return_value={"status": "ok", "result": [80, True]}), \ + patch.object(auditor.time, "sleep"): + yield run_worker + + +def _run(ctx): + # evaluate_lms(ctx, model_id, lmi, batch, submit, gi, packages_dir, no_cache) + evaluate_lms(ctx, 1, None, None, True, None, None, False) + + +def _read_cache(model_base_dir, model_index): + with open(_commit_store_path(model_base_dir, GI, BATCH, model_index)) as f: + return json.load(f) + + +def test_already_committed_lm_is_skipped_and_cache_untouched(tmp_path, worker_stubs): + """LM 0 committed on an earlier run, LM 1's commit failed: the retry must + leave LM 0's cache and on-chain commit alone and only commit LM 1.""" + state = {"committed": {0}, "commits": []} + ctx = _make_ctx(_make_auditor_contract(state), tmp_path) + _save_commit(tmp_path, GI, BATCH, 0, 55, False, OLD_SALT) + before = _read_cache(tmp_path, 0) + + with patch.object(auditor, "build_and_send_tx") as send: + _run(ctx) + + assert _read_cache(tmp_path, 0) == before + assert send.call_count == 1 + assert [args[:3] for args in state["commits"]] == [(GI, BATCH, 1)] + # Skipped before re-evaluating: the worker ran for LM 1 only. + assert worker_stubs.call_count == 1 + assert "lm-1" in worker_stubs.call_args.kwargs["container_name"] + + +def test_all_committed_sends_nothing(tmp_path, worker_stubs): + state = {"committed": set(MODELS), "commits": []} + ctx = _make_ctx(_make_auditor_contract(state), tmp_path) + + with patch.object(auditor, "build_and_send_tx") as send: + _run(ctx) + + send.assert_not_called() + worker_stubs.assert_not_called() + assert state["commits"] == [] + + +def test_cache_written_before_send_and_matches_sent_commit(tmp_path, worker_stubs): + """build_and_send_tx returning None (revert, or receipt wait failed on a + tx that may still mine) must still leave the preimage of the hash that + was actually sent on disk.""" + state = {"committed": set(), "commits": []} + ctx = _make_ctx(_make_auditor_contract(state), tmp_path) + seen_at_send = [] + + def _send(ctx_, fn, *args, **kwargs): + model_index = state["commits"][-1][2] + seen_at_send.append(_commit_store_path(tmp_path, GI, BATCH, model_index).exists()) + return None + + with patch.object(auditor, "build_and_send_tx", side_effect=_send): + _run(ctx) + + assert seen_at_send == [True, True] + assert len(state["commits"]) == 2 + for gi, batch_id, model_index, sent_hash in state["commits"]: + cached = _read_cache(tmp_path, model_index) + assert (cached["score"], cached["vote"]) == (80, True) + salt = bytes.fromhex(cached["salt"]) + assert _audit_commit_hash(80, True, salt, ACCOUNT, gi, batch_id, model_index) == sent_hash