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
2 changes: 1 addition & 1 deletion Documentation/public/roles/auditors.md
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,7 @@ dincli auditor lms-evaluation evaluate <model_id> [--gi <gi_index>] [--submit] [
| Option | Required | Description |
|---|---|---|
| `--gi <gi_index>` | 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 <batch_id>` | No | Evaluate a specific batch. Defaults to your assigned batch if omitted |
| `--lmi <lmi_index>` | No | Evaluate a single Local Model by index. Evaluates all models in the batch if omitted |

Expand Down
2 changes: 1 addition & 1 deletion Documentation/technical/contracts/DINTaskAuditor.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

---

Expand Down
2 changes: 1 addition & 1 deletion Documentation/technical/contracts/DINTaskCoordinator.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

---
Expand Down
11 changes: 7 additions & 4 deletions dincli/cli/aggregator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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}")

Expand Down
19 changes: 16 additions & 3 deletions dincli/cli/auditor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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),
Expand All @@ -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]")

Expand Down
105 changes: 105 additions & 0 deletions tests/test_aggregator_t2_batch_id.py
Original file line number Diff line number Diff line change
@@ -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)
187 changes: 187 additions & 0 deletions tests/test_auditor_commit_retry.py
Original file line number Diff line number Diff line change
@@ -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
Loading