iuna

iuna

iuna - experimental mainnet-candidate protocol
git clone https://getiuna.org/git/iuna.git
Log | Files | Refs | README | LICENSE

iuna_e2e.py (57419B)


      1 #!/usr/bin/env python3
      2 """Docker Compose e2e runner and reproducible chain checkpoint manager."""
      3 
      4 from __future__ import annotations
      5 
      6 import argparse
      7 from dataclasses import dataclass
      8 from datetime import datetime, timezone
      9 import hashlib
     10 import http.cookiejar
     11 import json
     12 import os
     13 from pathlib import Path
     14 import re
     15 import shutil
     16 import sqlite3
     17 import subprocess
     18 import sys
     19 import tempfile
     20 import time
     21 import urllib.parse
     22 import urllib.request
     23 import urllib.error
     24 
     25 
     26 ROOT = Path(__file__).resolve().parent.parent
     27 E2E_DIR = ROOT / "e2e"
     28 COMPOSE_FILES = (ROOT / "docker-compose.yml", E2E_DIR / "docker-compose.e2e.yml")
     29 SERVICES = ("bootstrap", "node2", "node3", "node4", "node5", "node6")
     30 SYNC_SERVICE = "syncnode"
     31 REFERENCE_MINING_ENV = "IUNA_E2E_REFERENCE_MINING_ENABLED"
     32 SERVICE_IPS = {
     33     "bootstrap": "172.29.0.10",
     34     "node2": "172.29.0.11",
     35     "node3": "172.29.0.12",
     36     "node4": "172.29.0.13",
     37     "node5": "172.29.0.14",
     38     "node6": "172.29.0.15",
     39 }
     40 PARTITION_GROUPS = (SERVICES[:3], SERVICES[3:])
     41 DEFAULT_PORTS = {
     42     "bootstrap": 28661,
     43     "node2": 28662,
     44     "node3": 28663,
     45     "node4": 28664,
     46     "node5": 28665,
     47     "node6": 28666,
     48     SYNC_SERVICE: 28667,
     49 }
     50 PORT_ENV = {
     51     "bootstrap": "IUNA_E2E_BOOTSTRAP_PORT",
     52     "node2": "IUNA_E2E_NODE2_PORT",
     53     "node3": "IUNA_E2E_NODE3_PORT",
     54     "node4": "IUNA_E2E_NODE4_PORT",
     55     "node5": "IUNA_E2E_NODE5_PORT",
     56     "node6": "IUNA_E2E_NODE6_PORT",
     57     SYNC_SERVICE: "IUNA_E2E_SYNCNODE_PORT",
     58 }
     59 SNAPSHOT_FILES = ("chain.sqlite3", "ui_data.sqlite3", "wallet.json", "config.json")
     60 EXPECTED_PROFILE = "iuna-local-e2e-5s-v1"
     61 EXPECTED_BLOCK_MS = 5_000
     62 HTTP_TIMEOUT_SECONDS = 30
     63 NAME_PATTERN = re.compile(r"^[a-z0-9][a-z0-9._-]*$")
     64 _OPENERS: dict[str, urllib.request.OpenerDirector] = {}
     65 
     66 
     67 @dataclass(frozen=True)
     68 class Scenario:
     69     snapshot: str
     70     through: int
     71     minimum_finalized_height: int | None = None
     72     canonical_height: int | None = None
     73     require_fixture_hash: bool = False
     74     leader_burn_minimum_height: int | None = None
     75 
     76 
     77 SCENARIOS = {
     78     "fallback-activation": Scenario(
     79         snapshot="pre-fallback-invalidation",
     80         through=300,
     81     ),
     82     "objective-finality": Scenario(
     83         snapshot="pre-objective-finality",
     84         through=1_001,
     85         minimum_finalized_height=1_000,
     86         canonical_height=1_000,
     87     ),
     88     "checkpoint-restart": Scenario(
     89         snapshot="first-objective-checkpoint",
     90         through=1_007,
     91         minimum_finalized_height=1_000,
     92         canonical_height=1_000,
     93         require_fixture_hash=True,
     94         leader_burn_minimum_height=1_002,
     95     ),
     96 }
     97 SPECIAL_SCENARIOS = ("sync-resilience", "partition-recovery")
     98 POST_ACTIVATION_SCENARIOS = (
     99     "objective-finality",
    100     "checkpoint-restart",
    101     *SPECIAL_SCENARIOS,
    102 )
    103 
    104 
    105 class E2EError(RuntimeError):
    106     pass
    107 
    108 
    109 def runtime_dir() -> Path:
    110     return Path(os.environ.get("IUNA_E2E_RUNTIME_DIR", E2E_DIR / ".runtime")).resolve()
    111 
    112 
    113 def snapshots_dir() -> Path:
    114     return Path(os.environ.get("IUNA_E2E_SNAPSHOTS_DIR", E2E_DIR / "snapshots")).resolve()
    115 
    116 
    117 def compose_env() -> dict[str, str]:
    118     env = os.environ.copy()
    119     env["IUNA_E2E_RUNTIME_DIR"] = str(runtime_dir())
    120     return env
    121 
    122 
    123 def compose_command(*args: str) -> list[str]:
    124     command = [
    125         "docker",
    126         "compose",
    127         "--project-name",
    128         os.environ.get("IUNA_E2E_PROJECT", "iuna-e2e"),
    129     ]
    130     for compose_file in COMPOSE_FILES:
    131         command.extend(("--file", str(compose_file)))
    132     command.extend(args)
    133     return command
    134 
    135 
    136 def compose(*args: str, check: bool = True) -> subprocess.CompletedProcess[str]:
    137     return subprocess.run(
    138         compose_command(*args),
    139         cwd=ROOT,
    140         env=compose_env(),
    141         check=check,
    142         text=True,
    143     )
    144 
    145 
    146 def create_evidence_run(base: Path | None, scenario: str) -> tuple[Path | None, dict]:
    147     started_at = datetime.now(timezone.utc)
    148     report = {
    149         "format": 1,
    150         "scenario": scenario,
    151         "started_at": started_at.isoformat(),
    152         "git_commit": subprocess.run(
    153             ["git", "rev-parse", "HEAD"],
    154             cwd=ROOT,
    155             check=True,
    156             capture_output=True,
    157             text=True,
    158         ).stdout.strip(),
    159         "git_dirty": bool(
    160             subprocess.run(
    161                 ["git", "status", "--porcelain"],
    162                 cwd=ROOT,
    163                 check=True,
    164                 capture_output=True,
    165                 text=True,
    166             ).stdout.strip()
    167         ),
    168         "tracked_tree_sha256": tracked_tree_sha256(),
    169         "phases": {},
    170         "outcome": "running",
    171     }
    172     if base is None:
    173         return None, report
    174     run = base.resolve() / f"{started_at.strftime('%Y%m%dT%H%M%S.%fZ')}-{scenario}"
    175     run.mkdir(parents=True, exist_ok=False)
    176     write_evidence_report(run, report)
    177     return run, report
    178 
    179 
    180 def tracked_tree_sha256() -> str:
    181     tracked = subprocess.run(
    182         ["git", "ls-files", "-z"],
    183         cwd=ROOT,
    184         check=True,
    185         capture_output=True,
    186     ).stdout.split(b"\0")
    187     digest = hashlib.sha256()
    188     for encoded_path in tracked:
    189         if not encoded_path:
    190             continue
    191         path = ROOT / os.fsdecode(encoded_path)
    192         digest.update(encoded_path)
    193         digest.update(b"\0")
    194         if path.exists():
    195             contents = path.read_bytes()
    196             digest.update(len(contents).to_bytes(8, "big"))
    197             digest.update(contents)
    198         else:
    199             digest.update(b"missing")
    200     return digest.hexdigest()
    201 
    202 
    203 def write_evidence_report(run: Path | None, report: dict) -> None:
    204     if run is None:
    205         return
    206     temporary = run / "report.json.tmp"
    207     temporary.write_text(json.dumps(report, indent=2, sort_keys=True) + "\n")
    208     temporary.replace(run / "report.json")
    209 
    210 
    211 def capture_evidence_logs(
    212     run: Path | None, services: tuple[str, ...] = SERVICES
    213 ) -> None:
    214     if run is None:
    215         return
    216     result = subprocess.run(
    217         compose_command("logs", "--no-color", *services),
    218         cwd=ROOT,
    219         env=compose_env(),
    220         check=False,
    221         capture_output=True,
    222         text=True,
    223     )
    224     (run / "nodes.log").write_text(result.stdout + result.stderr)
    225 
    226 
    227 def evidence_block(block: dict) -> dict:
    228     return {
    229         key: block.get(key)
    230         for key in (
    231             "height",
    232             "hash",
    233             "prev_hash",
    234             "timestamp_ms",
    235             "finalizer_mode",
    236             "finalizer_rank",
    237             "miner",
    238         )
    239     }
    240 
    241 
    242 def evidence_statuses(statuses: dict[str, dict]) -> dict[str, dict]:
    243     return {service: compact_status(status) for service, status in statuses.items()}
    244 
    245 
    246 def container_command(
    247     service: str, *args: str, check: bool = True
    248 ) -> subprocess.CompletedProcess[str]:
    249     return subprocess.run(
    250         compose_command("exec", "-T", service, *args),
    251         cwd=ROOT,
    252         env=compose_env(),
    253         check=check,
    254         text=True,
    255         stdout=None if check else subprocess.DEVNULL,
    256         stderr=None if check else subprocess.DEVNULL,
    257     )
    258 
    259 
    260 def clear_partition() -> None:
    261     for service in SERVICES:
    262         for builtin, chain in (
    263             ("INPUT", "IUNA_E2E_INPUT"),
    264             ("OUTPUT", "IUNA_E2E_OUTPUT"),
    265         ):
    266             container_command(
    267                 service, "iptables", "-D", builtin, "-j", chain, check=False
    268             )
    269             container_command(service, "iptables", "-F", chain, check=False)
    270             container_command(service, "iptables", "-X", chain, check=False)
    271 
    272 
    273 def apply_partition() -> None:
    274     clear_partition()
    275     left, right = PARTITION_GROUPS
    276     for group, blocked_group in ((left, right), (right, left)):
    277         for service in group:
    278             for builtin, chain in (
    279                 ("INPUT", "IUNA_E2E_INPUT"),
    280                 ("OUTPUT", "IUNA_E2E_OUTPUT"),
    281             ):
    282                 container_command(service, "iptables", "-N", chain)
    283                 container_command(service, "iptables", "-I", builtin, "1", "-j", chain)
    284             for blocked_service in blocked_group:
    285                 blocked_ip = SERVICE_IPS[blocked_service]
    286                 container_command(
    287                     service,
    288                     "iptables",
    289                     "-A",
    290                     "IUNA_E2E_INPUT",
    291                     "-s",
    292                     blocked_ip,
    293                     "-j",
    294                     "REJECT",
    295                 )
    296                 container_command(
    297                     service,
    298                     "iptables",
    299                     "-A",
    300                     "IUNA_E2E_OUTPUT",
    301                     "-d",
    302                     blocked_ip,
    303                     "-j",
    304                     "REJECT",
    305                 )
    306     print(f"partition active: {left} | {right}", flush=True)
    307 
    308 
    309 def validate_snapshot_name(name: str) -> str:
    310     if not NAME_PATTERN.fullmatch(name):
    311         raise E2EError(
    312             "snapshot names must start with a lowercase letter or digit and only contain "
    313             "lowercase letters, digits, '.', '_' or '-'"
    314         )
    315     return name
    316 
    317 
    318 def snapshot_path(name: str) -> Path:
    319     return snapshots_dir() / validate_snapshot_name(name)
    320 
    321 
    322 def node_port(service: str) -> int:
    323     return int(os.environ.get(PORT_ENV[service], DEFAULT_PORTS[service]))
    324 
    325 
    326 def node_json(service: str, path: str) -> object:
    327     base_url = f"http://127.0.0.1:{node_port(service)}"
    328     opener = _OPENERS.get(service)
    329     if opener is None:
    330         password = os.environ.get("IUNA_TESTNET_PASSWORD", "testtesttest")
    331         jar = http.cookiejar.CookieJar()
    332         opener = urllib.request.build_opener(urllib.request.HTTPCookieProcessor(jar))
    333         login = urllib.request.Request(
    334             f"{base_url}/api/auth/login",
    335             data=urllib.parse.urlencode({"password": password}).encode(),
    336             headers={
    337                 "Content-Type": "application/x-www-form-urlencoded",
    338                 "Origin": base_url,
    339             },
    340             method="POST",
    341         )
    342         with opener.open(login, timeout=HTTP_TIMEOUT_SECONDS) as response:
    343             result = json.load(response)
    344         if not result.get("ok"):
    345             raise E2EError(f"{service} login failed: {result.get('error', 'unknown error')}")
    346         _OPENERS[service] = opener
    347     with opener.open(
    348         f"{base_url}{path}", timeout=HTTP_TIMEOUT_SECONDS
    349     ) as response:
    350         return json.load(response)
    351 
    352 
    353 def node_form(service: str, path: str, values: dict[str, object]) -> dict:
    354     node_json(service, "/api/status")
    355     base_url = f"http://127.0.0.1:{node_port(service)}"
    356     request = urllib.request.Request(
    357         f"{base_url}{path}",
    358         data=urllib.parse.urlencode(values).encode(),
    359         headers={"Content-Type": "application/x-www-form-urlencoded", "Origin": base_url},
    360         method="POST",
    361     )
    362     with _OPENERS[service].open(request, timeout=HTTP_TIMEOUT_SECONDS) as response:
    363         result = json.load(response)
    364     if not isinstance(result, dict) or not result.get("ok"):
    365         raise E2EError(f"{service} form request failed for {path}: {result}")
    366     return result
    367 
    368 
    369 def node_status(service: str) -> dict:
    370     result = node_json(service, "/api/status")
    371     if not isinstance(result, dict):
    372         raise E2EError(f"{service} returned a non-object status response")
    373     return result
    374 
    375 
    376 def all_statuses() -> dict[str, dict]:
    377     return {service: node_status(service) for service in SERVICES}
    378 
    379 
    380 def compact_status(status: dict) -> dict:
    381     return {
    382         "height": status["chain"]["height"],
    383         "tip_hash": status["chain"]["tip_hash"],
    384         "finalized_height": status["chain"].get("finalized_height"),
    385         "profile": status["launch_profile"]["profile_id"],
    386         "target_block_ms": status["mining"]["vdf_target_block_ms"],
    387     }
    388 
    389 
    390 def assert_e2e_profile(statuses: dict[str, dict]) -> None:
    391     errors = []
    392     for service, status in statuses.items():
    393         profile = status["launch_profile"]["profile_id"]
    394         block_ms = status["mining"]["vdf_target_block_ms"]
    395         if profile != EXPECTED_PROFILE or block_ms != EXPECTED_BLOCK_MS:
    396             errors.append(f"{service}: profile={profile}, target_block_ms={block_ms}")
    397     if errors:
    398         raise E2EError("not running the isolated 5s e2e build: " + "; ".join(errors))
    399 
    400 
    401 def assert_converged(statuses: dict[str, dict]) -> None:
    402     tips = {
    403         (status["chain"]["height"], status["chain"]["tip_hash"])
    404         for status in statuses.values()
    405     }
    406     if len(tips) != 1:
    407         summary = ", ".join(
    408             f"{service}={status['chain']['height']}:{status['chain']['tip_hash'][:12]}"
    409             for service, status in statuses.items()
    410         )
    411         raise E2EError(f"nodes did not converge: {summary}")
    412 
    413 
    414 def assert_finality(statuses: dict[str, dict], minimum_height: int | None) -> None:
    415     checkpoints = {
    416         (status["chain"].get("finalized_height"), status["chain"].get("finalized_hash"))
    417         for status in statuses.values()
    418     }
    419     if len(checkpoints) != 1:
    420         raise E2EError(f"nodes disagree about objective finality: {checkpoints}")
    421     finalized_height, finalized_hash = next(iter(checkpoints))
    422     if minimum_height is None:
    423         if finalized_height is not None or finalized_hash is not None:
    424             raise E2EError(
    425                 f"objective finality activated too early at height {finalized_height}"
    426             )
    427         return
    428     if finalized_height is None or finalized_height < minimum_height or not finalized_hash:
    429         raise E2EError(
    430             f"expected a finalized checkpoint at or above {minimum_height}, "
    431             f"got {finalized_height}:{finalized_hash}"
    432         )
    433     block_hashes = {
    434         block_at_height(service, finalized_height)["hash"] for service in SERVICES
    435     }
    436     if block_hashes != {finalized_hash}:
    437         raise E2EError(
    438             f"finalized hash does not identify block {finalized_height}: {block_hashes}"
    439         )
    440 
    441 
    442 def block_at_height(service: str, height: int) -> dict:
    443     result = node_json(service, f"/api/blocks?before_height={height + 1}&limit=1")
    444     if not isinstance(result, list) or len(result) != 1 or result[0].get("height") != height:
    445         raise E2EError(f"{service} block API did not return block {height}")
    446     return result[0]
    447 
    448 
    449 def recent_blocks(service: str, limit: int = 100) -> list[dict]:
    450     result = node_json(service, f"/api/blocks?limit={limit}")
    451     if not isinstance(result, list):
    452         raise E2EError(f"{service} block API returned a non-list response")
    453     return result
    454 
    455 
    456 def wait_for_partition_recovery(
    457     after_heights: dict[str, int], timeout: float
    458 ) -> tuple[dict[str, dict], dict[str, int]]:
    459     deadline = time.monotonic() + timeout
    460     last_summary = "nodes unavailable"
    461     left, right = PARTITION_GROUPS
    462     while time.monotonic() < deadline:
    463         try:
    464             statuses = all_statuses()
    465             assert_e2e_profile(statuses)
    466             group_tips = []
    467             recovery_heights: dict[str, int] = {}
    468             for label, group in (("left", left), ("right", right)):
    469                 tips = {
    470                     (
    471                         statuses[service]["chain"]["height"],
    472                         statuses[service]["chain"]["tip_hash"],
    473                     )
    474                     for service in group
    475                 }
    476                 if len(tips) != 1:
    477                     break
    478                 group_tip = next(iter(tips))
    479                 group_tips.append(group_tip)
    480                 recoveries = [
    481                     int(block["height"])
    482                     for block in recent_blocks(group[0])
    483                     if block.get("finalizer_mode") == "recovery"
    484                     and int(block.get("height", -1)) > after_heights[label]
    485                 ]
    486                 if recoveries:
    487                     recovery_heights[label] = max(recoveries)
    488             last_summary = ", ".join(
    489                 f"{service}={status['chain']['height']}:{status['chain']['tip_hash'][:12]}"
    490                 for service, status in statuses.items()
    491             )
    492             if (
    493                 len(group_tips) == 2
    494                 and group_tips[0] != group_tips[1]
    495                 and set(recovery_heights) == {"left", "right"}
    496             ):
    497                 return statuses, recovery_heights
    498         except (
    499             E2EError,
    500             OSError,
    501             KeyError,
    502             ValueError,
    503             urllib.error.URLError,
    504         ) as error:
    505             last_summary = str(error)
    506         time.sleep(0.25)
    507     raise E2EError(
    508         "timed out waiting for divergent recovery blocks in both partitions; "
    509         + last_summary
    510     )
    511 
    512 
    513 def automatic_finalization_settings(statuses: dict[str, dict]) -> dict[str, dict]:
    514     return {
    515         service: {
    516             "enabled": bool(statuses[service]["mining"]["automatic"]),
    517             "amount": int(statuses[service]["mining"]["burn_per_block"]),
    518             "fee_per_byte": int(statuses[service]["mining"]["automatic_burn_fee"]),
    519         }
    520         for service in SERVICES
    521     }
    522 
    523 
    524 def configure_partition_recovery_workers(
    525     statuses: dict[str, dict], settings: dict[str, dict], timeout: float
    526 ) -> dict[str, str]:
    527     workers = {
    528         "left": PARTITION_GROUPS[0][0],
    529         "right": PARTITION_GROUPS[1][0],
    530     }
    531     for label, group in (("left", PARTITION_GROUPS[0]), ("right", PARTITION_GROUPS[1])):
    532         tips = {statuses[service]["chain"]["tip_hash"] for service in group}
    533         if len(tips) != 1:
    534             raise E2EError(f"cannot select {label} recovery worker before convergence")
    535     if not all(values["enabled"] for values in settings.values()):
    536         raise E2EError("partition recovery requires automatic finalization on every node")
    537 
    538     for service in SERVICES:
    539         mining = statuses[service]["mining"]
    540         if service == workers["left"]:
    541             continue
    542         node_form(
    543             service,
    544             "/api/settings/burn-per-block",
    545             {
    546                 "enabled": "false",
    547                 "amount": mining["burn_per_block"],
    548                 "fee_per_byte": mining["automatic_burn_fee"],
    549             },
    550         )
    551 
    552     # The automatic-finalizer loop polls this setting and cancels any in-flight
    553     # VDF. Let every disabled node observe it, then absorb a block that may have
    554     # crossed the publication boundary concurrently with the settings update.
    555     time.sleep(2)
    556     current = all_statuses()
    557     target = max(status["chain"]["height"] for status in current.values())
    558     wait_for_height(target, timeout, converge=True)
    559 
    560     right_settings = settings[workers["right"]]
    561     node_form(
    562         workers["right"],
    563         "/api/settings/burn-per-block",
    564         {
    565             "enabled": "true",
    566             "amount": right_settings["amount"],
    567             "fee_per_byte": right_settings["fee_per_byte"],
    568         },
    569     )
    570     print(
    571         "partition recovery workers: " + json.dumps(workers, sort_keys=True),
    572         flush=True,
    573     )
    574     return workers
    575 
    576 
    577 def restore_automatic_finalization(settings: dict[str, dict]) -> None:
    578     for service, values in settings.items():
    579         node_form(
    580             service,
    581             "/api/settings/burn-per-block",
    582             {
    583                 "enabled": "true" if values["enabled"] else "false",
    584                 "amount": values["amount"],
    585                 "fee_per_byte": values["fee_per_byte"],
    586             },
    587         )
    588 
    589 
    590 def wait_for_ticket_after(height: int, timeout: float) -> tuple[dict[str, dict], dict]:
    591     deadline = time.monotonic() + timeout
    592     last_summary = "nodes unavailable"
    593     while time.monotonic() < deadline:
    594         try:
    595             statuses = all_statuses()
    596             assert_e2e_profile(statuses)
    597             assert_converged(statuses)
    598             tickets = [
    599                 block
    600                 for block in recent_blocks(SERVICES[0])
    601                 if block.get("finalizer_mode") == "ticket"
    602                 and int(block.get("height", -1)) > height
    603                 and int(block.get("finalizer_rank", -1)) == 0
    604             ]
    605             if tickets:
    606                 ticket = min(tickets, key=lambda block: int(block["height"]))
    607                 finalized_height = statuses[SERVICES[0]]["chain"].get(
    608                     "finalized_height"
    609                 )
    610                 if finalized_height is not None and finalized_height >= height:
    611                     return statuses, ticket
    612             last_summary = ", ".join(
    613                 f"{service}={status['chain']['height']}"
    614                 for service, status in statuses.items()
    615             )
    616         except (
    617             E2EError,
    618             OSError,
    619             KeyError,
    620             ValueError,
    621             urllib.error.URLError,
    622         ) as error:
    623             last_summary = str(error)
    624         time.sleep(0.25)
    625     raise E2EError(
    626         f"timed out waiting for a rank-0 ticket after {height}; {last_summary}"
    627     )
    628 
    629 
    630 def assert_leader_uses_burn_from_height(
    631     block_height: int, minimum_height: int
    632 ) -> None:
    633     block = block_at_height(SERVICES[0], block_height)
    634     leader_proof = block.get("leader_proof")
    635     if block.get("finalizer_mode") != "ticket" or not isinstance(leader_proof, dict):
    636         raise E2EError(f"block {block_height} was not finalized by a burn ticket")
    637 
    638     ticket_id = leader_proof.get("ticket_id")
    639     if not isinstance(ticket_id, str) or not ticket_id:
    640         raise E2EError(f"block {block_height} has no valid leader ticket ID")
    641 
    642     for height in range(minimum_height, block_height):
    643         source = block_at_height(SERVICES[0], height)
    644         if any(
    645             transaction.get("kind") == "burn"
    646             and transaction.get("signature") == ticket_id
    647             for transaction in source.get("transactions", [])
    648         ):
    649             print(
    650                 f"block {block_height} leader ticket comes from burn at height {height}"
    651             )
    652             return
    653 
    654     raise E2EError(
    655         f"block {block_height} leader ticket does not come from a burn at or after "
    656         f"height {minimum_height}"
    657     )
    658 
    659 
    660 def assert_api_health(
    661     statuses: dict[str, dict], through: int, services: tuple[str, ...] = SERVICES
    662 ) -> None:
    663     for service in services:
    664         blocks = node_json(service, "/api/blocks?limit=3")
    665         health = node_json(service, "/api/network/health")
    666         peers = node_json(service, "/api/peers?limit=100")
    667         mempool = node_json(service, "/api/mempool?limit=100")
    668         wallet = node_json(service, "/api/wallet/transactions?limit=3")
    669         if not isinstance(blocks, list) or not blocks:
    670             raise E2EError(f"{service} block API returned no blocks")
    671         if max(block.get("height", -1) for block in blocks) < through:
    672             raise E2EError(f"{service} block API has not projected height {through}")
    673         if not isinstance(health, dict) or health.get("local_height", -1) < through:
    674             raise E2EError(f"{service} network health is behind height {through}")
    675         for label, page in (("peers", peers), ("mempool", mempool), ("wallet", wallet)):
    676             if not isinstance(page, dict) or not isinstance(page.get("items"), list):
    677                 raise E2EError(f"{service} {label} API returned an invalid page")
    678         if health.get("local_tip_hash") != statuses[service]["chain"]["tip_hash"]:
    679             # The chain is live; advancing between the status and health requests is valid.
    680             if health.get("local_height", 0) <= statuses[service]["chain"]["height"]:
    681                 raise E2EError(f"{service} status and network health disagree about the tip")
    682 
    683 
    684 def read_chain_metadata(path: Path) -> dict:
    685     if not path.exists():
    686         raise E2EError(f"chain database does not exist: {path}")
    687     connection = sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=2)
    688     try:
    689         row = connection.execute(
    690             "SELECT height, tip_hash, updated_at_ms FROM chain_snapshots WHERE id = 1"
    691         ).fetchone()
    692     finally:
    693         connection.close()
    694     if row is None:
    695         raise E2EError(f"chain database has no persisted snapshot: {path}")
    696     return {"height": row[0], "tip_hash": row[1], "updated_at_ms": row[2]}
    697 
    698 
    699 def remove_node_databases(service: str) -> None:
    700     directory = runtime_dir() / service
    701     for filename in ("chain.sqlite3", "ui_data.sqlite3"):
    702         database = directory / filename
    703         for suffix in ("", "-wal", "-shm", "-journal"):
    704             candidate = Path(f"{database}{suffix}")
    705             if candidate.exists():
    706                 candidate.unlink()
    707 
    708 
    709 def restore_node_databases(
    710     service: str, snapshot: str, fixture_service: str | None = None
    711 ) -> dict:
    712     source = snapshot_path(snapshot)
    713     verify_snapshot(source)
    714     target = runtime_dir() / service
    715     source_service = fixture_service or service
    716     remove_node_databases(service)
    717     for filename in ("chain.sqlite3", "ui_data.sqlite3"):
    718         shutil.copy2(source / source_service / filename, target / filename)
    719     return read_chain_metadata(target / "chain.sqlite3")
    720 
    721 
    722 def wait_for_active_sync(service: str, minimum_target: int, timeout: float) -> dict:
    723     deadline = time.monotonic() + timeout
    724     last_summary = "node unavailable"
    725     while time.monotonic() < deadline:
    726         try:
    727             health = node_json(service, "/api/network/health")
    728             start = health.get("sync_start_height")
    729             validated = health.get("sync_validated_height")
    730             target = health.get("sync_target_height")
    731             last_summary = (
    732                 f"local={health.get('local_height')}, start={start}, "
    733                 f"validated={validated}, target={target}"
    734             )
    735             if (
    736                 isinstance(start, int)
    737                 and isinstance(validated, int)
    738                 and isinstance(target, int)
    739                 and target >= minimum_target
    740                 and validated < target
    741             ):
    742                 return {
    743                     "local_height": health.get("local_height"),
    744                     "start_height": start,
    745                     "validated_height": validated,
    746                     "target_height": target,
    747                 }
    748         except (
    749             E2EError,
    750             OSError,
    751             KeyError,
    752             ValueError,
    753             urllib.error.URLError,
    754         ) as error:
    755             last_summary = str(error)
    756         time.sleep(0.01)
    757     raise E2EError(
    758         f"timed out waiting for active range sync on {service}; {last_summary}"
    759     )
    760 
    761 
    762 def wait_for_persisted_height(target: int, timeout: float) -> dict:
    763     deadline = time.monotonic() + timeout
    764     database = runtime_dir() / "bootstrap" / "chain.sqlite3"
    765     last_height = None
    766     while time.monotonic() < deadline:
    767         try:
    768             metadata = read_chain_metadata(database)
    769             last_height = metadata["height"]
    770             if last_height >= target:
    771                 return metadata
    772         except sqlite3.Error:
    773             pass
    774         time.sleep(0.1)
    775     raise E2EError(
    776         f"timed out waiting for persisted bootstrap height {target}; last height={last_height}"
    777     )
    778 
    779 
    780 def materialize_chain_checkpoint(
    781     source: Path, chain_destination: Path, ui_destination: Path, height: int
    782 ) -> None:
    783     subprocess.run(
    784         [
    785             "cargo",
    786             "run",
    787             "--quiet",
    788             "--features",
    789             "e2e",
    790             "--bin",
    791             "iuna-e2e-checkpoint",
    792             "--",
    793             str(source),
    794             str(chain_destination),
    795             str(ui_destination),
    796             str(height),
    797         ],
    798         cwd=ROOT,
    799         check=True,
    800         text=True,
    801     )
    802 
    803 
    804 def wait_for_height(
    805     target: int,
    806     timeout: float,
    807     converge: bool,
    808     services: tuple[str, ...] = SERVICES,
    809 ) -> dict[str, dict]:
    810     deadline = time.monotonic() + timeout
    811     last_summary = "nodes unavailable"
    812     while time.monotonic() < deadline:
    813         try:
    814             statuses = {service: node_status(service) for service in services}
    815             assert_e2e_profile(statuses)
    816             tips = {
    817                 (status["chain"]["height"], status["chain"]["tip_hash"])
    818                 for status in statuses.values()
    819             }
    820             heights = [status["chain"]["height"] for status in statuses.values()]
    821             last_summary = ", ".join(
    822                 f"{service}={status['chain']['height']}" for service, status in statuses.items()
    823             )
    824             reached = min(heights) >= target
    825             if reached and (not converge or len(tips) == 1):
    826                 return statuses
    827         except (E2EError, OSError, KeyError, urllib.error.URLError) as error:
    828             last_summary = str(error)
    829         time.sleep(0.25)
    830     raise E2EError(f"timed out waiting for height {target}: {last_summary}")
    831 
    832 
    833 def sqlite_backup(source: Path, destination: Path) -> None:
    834     source_connection = sqlite3.connect(f"file:{source}?mode=ro", uri=True, timeout=5)
    835     destination_connection = sqlite3.connect(destination)
    836     try:
    837         source_connection.backup(destination_connection)
    838         destination_connection.execute("PRAGMA journal_mode = DELETE")
    839     finally:
    840         destination_connection.close()
    841         source_connection.close()
    842 
    843 
    844 def sha256(path: Path) -> str:
    845     digest = hashlib.sha256()
    846     with path.open("rb") as source:
    847         for chunk in iter(lambda: source.read(1024 * 1024), b""):
    848             digest.update(chunk)
    849     return digest.hexdigest()
    850 
    851 
    852 def snapshot_checksums(directory: Path) -> dict[str, str]:
    853     return {
    854         str(path.relative_to(directory)): sha256(path)
    855         for path in sorted(directory.glob("*/*"))
    856         if path.is_file()
    857     }
    858 
    859 
    860 def verify_snapshot(directory: Path) -> dict:
    861     manifest_path = directory / "manifest.json"
    862     if not manifest_path.exists():
    863         raise E2EError(f"snapshot has no manifest: {directory}")
    864     manifest = json.loads(manifest_path.read_text())
    865     expected = manifest.get("sha256", {})
    866     actual = snapshot_checksums(directory)
    867     if actual != expected:
    868         raise E2EError(f"snapshot checksum mismatch: {directory}")
    869     return manifest
    870 
    871 
    872 def test_snapshots(plan_path: Path = E2E_DIR / "checkpoints.json") -> None:
    873     plan = json.loads(plan_path.read_text())
    874     expected_names = {checkpoint["name"] for checkpoint in plan}
    875     scenario_snapshots = {scenario.snapshot for scenario in SCENARIOS.values()}
    876     missing_scenarios = scenario_snapshots - expected_names
    877     if missing_scenarios:
    878         raise E2EError(f"scenarios reference snapshots outside the plan: {missing_scenarios}")
    879 
    880     for checkpoint in plan:
    881         name = validate_snapshot_name(checkpoint["name"])
    882         expected_height = int(checkpoint["height"])
    883         directory = snapshot_path(name)
    884         manifest = verify_snapshot(directory)
    885         if manifest.get("name") != name or manifest.get("height") != expected_height:
    886             raise E2EError(
    887                 f"{name} manifest identity mismatch: "
    888                 f"{manifest.get('name')} at {manifest.get('height')}"
    889             )
    890         if manifest.get("profile_id") != EXPECTED_PROFILE:
    891             raise E2EError(f"{name} uses profile {manifest.get('profile_id')}")
    892         if manifest.get("target_block_ms") != EXPECTED_BLOCK_MS:
    893             raise E2EError(f"{name} uses target {manifest.get('target_block_ms')}ms")
    894         nodes = manifest.get("nodes", {})
    895         if set(nodes) != set(SERVICES):
    896             raise E2EError(f"{name} does not contain exactly the six e2e nodes")
    897         for service in SERVICES:
    898             metadata = read_chain_metadata(directory / service / "chain.sqlite3")
    899             recorded = nodes[service]
    900             if metadata != recorded:
    901                 raise E2EError(f"{name}/{service} chain metadata differs from its manifest")
    902             if metadata["height"] != expected_height:
    903                 raise E2EError(
    904                     f"{name}/{service} is at {metadata['height']}, expected {expected_height}"
    905                 )
    906             if metadata["tip_hash"] != manifest.get("tip_hash"):
    907                 raise E2EError(f"{name}/{service} has a different canonical tip")
    908     print(f"e2e snapshots passed ({len(plan)} checkpoints)")
    909 
    910 
    911 def assert_canonical_block(height: int, require_fixture_hash: bool) -> None:
    912     hashes = {block_at_height(service, height)["hash"] for service in SERVICES}
    913     if len(hashes) != 1:
    914         raise E2EError(f"nodes disagree about canonical block {height}: {hashes}")
    915     if not require_fixture_hash:
    916         return
    917     checkpoint = next(
    918         (
    919             item
    920             for item in json.loads((E2E_DIR / "checkpoints.json").read_text())
    921             if int(item["height"]) == height
    922         ),
    923         None,
    924     )
    925     if checkpoint is None:
    926         raise E2EError(f"no checkpoint fixture records canonical block {height}")
    927     expected_hash = verify_snapshot(snapshot_path(checkpoint["name"]))["tip_hash"]
    928     if hashes != {expected_hash}:
    929         raise E2EError(
    930             f"canonical block {height} differs from the checkpoint: {hashes}"
    931         )
    932 
    933 
    934 def capture_snapshot(name: str, height: int, timeout: float, force: bool) -> None:
    935     destination = snapshot_path(name)
    936     if destination.exists() and not force:
    937         raise E2EError(f"snapshot already exists: {destination}; pass --force to replace it")
    938     assert_e2e_profile({"bootstrap": node_status("bootstrap")})
    939     live_metadata = wait_for_persisted_height(height, timeout)
    940     compose("pause", *SERVICES)
    941     temp_path: Path | None = None
    942     try:
    943         snapshots_dir().mkdir(parents=True, exist_ok=True)
    944         temp_path = Path(tempfile.mkdtemp(prefix=f".{name}-", dir=snapshots_dir()))
    945         canonical_database = temp_path / ".canonical-chain.sqlite3"
    946         canonical_ui_database = temp_path / ".canonical-ui-data.sqlite3"
    947         materialize_chain_checkpoint(
    948             runtime_dir() / "bootstrap" / "chain.sqlite3",
    949             canonical_database,
    950             canonical_ui_database,
    951             height,
    952         )
    953         nodes = {}
    954         for service in SERVICES:
    955             source_dir = runtime_dir() / service
    956             target_dir = temp_path / service
    957             target_dir.mkdir()
    958             sqlite_backup(canonical_database, target_dir / "chain.sqlite3")
    959             sqlite_backup(canonical_ui_database, target_dir / "ui_data.sqlite3")
    960             for filename in SNAPSHOT_FILES[2:]:
    961                 source = source_dir / filename
    962                 if not source.exists():
    963                     raise E2EError(f"required {service} snapshot file is missing: {source}")
    964                 shutil.copy2(source, target_dir / filename)
    965             nodes[service] = read_chain_metadata(target_dir / "chain.sqlite3")
    966 
    967         for database in (canonical_database, canonical_ui_database):
    968             for suffix in ("", "-shm", "-wal"):
    969                 candidate = Path(f"{database}{suffix}")
    970                 if candidate.exists():
    971                     candidate.unlink()
    972 
    973         canonical = nodes["bootstrap"]
    974         manifest = {
    975             "format": 1,
    976             "name": name,
    977             "height": canonical["height"],
    978             "tip_hash": canonical["tip_hash"],
    979             "profile_id": EXPECTED_PROFILE,
    980             "target_block_ms": EXPECTED_BLOCK_MS,
    981             "source_height": live_metadata["height"],
    982             "nodes": nodes,
    983             "sha256": snapshot_checksums(temp_path),
    984         }
    985         (temp_path / "manifest.json").write_text(json.dumps(manifest, indent=2) + "\n")
    986         if destination.exists():
    987             shutil.rmtree(destination)
    988         temp_path.rename(destination)
    989         temp_path = None
    990         print(f"captured {name} at height {height}: {destination}")
    991     finally:
    992         if temp_path is not None:
    993             shutil.rmtree(temp_path, ignore_errors=True)
    994         compose("unpause", *SERVICES, check=False)
    995 
    996 
    997 def restore_snapshot(name: str) -> None:
    998     source = snapshot_path(name)
    999     manifest = verify_snapshot(source)
   1000     if manifest.get("profile_id") != EXPECTED_PROFILE:
   1001         raise E2EError(
   1002             f"snapshot profile is {manifest.get('profile_id')}, expected {EXPECTED_PROFILE}"
   1003         )
   1004     _OPENERS.clear()
   1005     compose("down", "--remove-orphans", check=False)
   1006     base = runtime_dir()
   1007     base.mkdir(parents=True, exist_ok=True)
   1008     for service in SERVICES:
   1009         target = base / service
   1010         target.mkdir(parents=True, exist_ok=True)
   1011         for existing in target.iterdir():
   1012             if existing.is_dir() and not existing.is_symlink():
   1013                 shutil.rmtree(existing)
   1014             else:
   1015                 existing.unlink()
   1016         for filename in SNAPSHOT_FILES:
   1017             snapshot_file = source / service / filename
   1018             if filename == "ui_data.sqlite3" and not snapshot_file.exists():
   1019                 continue
   1020             shutil.copy2(snapshot_file, target / filename)
   1021     print(f"restored {name} at height {manifest['height']} into {base}")
   1022 
   1023 
   1024 def reset_runtime() -> None:
   1025     compose("down", "--remove-orphans", check=False)
   1026     base = runtime_dir()
   1027     for service in (*SERVICES, SYNC_SERVICE):
   1028         target = base / service
   1029         if target.exists():
   1030             shutil.rmtree(target)
   1031     _OPENERS.clear()
   1032     print(f"cleared mutable e2e node data under {base}")
   1033 
   1034 
   1035 def start(build: bool) -> None:
   1036     _OPENERS.clear()
   1037     runtime_dir().mkdir(parents=True, exist_ok=True)
   1038     for service in SERVICES:
   1039         (runtime_dir() / service).mkdir(parents=True, exist_ok=True)
   1040     args = ["up", "--detach", "--wait", "--wait-timeout", "600"]
   1041     if build:
   1042         args.append("--build")
   1043     args.extend(SERVICES)
   1044     compose(*args)
   1045     statuses = all_statuses()
   1046     assert_e2e_profile(statuses)
   1047     print(json.dumps({name: compact_status(value) for name, value in statuses.items()}, indent=2))
   1048 
   1049 
   1050 def capture_plan(plan_path: Path, timeout: float, force: bool) -> None:
   1051     plan = json.loads(plan_path.read_text())
   1052     for checkpoint in plan:
   1053         name = validate_snapshot_name(checkpoint["name"])
   1054         height = int(checkpoint["height"])
   1055         destination = snapshot_path(name)
   1056         if destination.exists() and not force:
   1057             print(f"keeping existing snapshot {name}")
   1058             continue
   1059         capture_snapshot(name, height, timeout, force)
   1060 
   1061 
   1062 def smoke(name: str, through: int, timeout: float, build: bool, keep: bool) -> None:
   1063     restore_snapshot(name)
   1064     try:
   1065         start(build)
   1066         statuses = wait_for_height(through, timeout, converge=True)
   1067         summary = {service: compact_status(status) for service, status in statuses.items()}
   1068         print(f"e2e smoke passed through height {through}")
   1069         print(json.dumps(summary, indent=2))
   1070     except Exception:
   1071         compose("logs", "--tail", "200", *SERVICES, check=False)
   1072         raise
   1073     finally:
   1074         if not keep:
   1075             compose("down", "--remove-orphans", check=False)
   1076 
   1077 
   1078 def run_scenario(
   1079     name: str, scenario: Scenario, timeout: float, build: bool, keep: bool
   1080 ) -> None:
   1081     print(
   1082         f"running e2e scenario {name}: {scenario.snapshot} -> {scenario.through}",
   1083         flush=True,
   1084     )
   1085     restore_snapshot(scenario.snapshot)
   1086     try:
   1087         start(build)
   1088         statuses = wait_for_height(scenario.through, timeout, converge=True)
   1089         assert_converged(statuses)
   1090         assert_finality(statuses, scenario.minimum_finalized_height)
   1091         if scenario.canonical_height is not None:
   1092             assert_canonical_block(
   1093                 scenario.canonical_height, scenario.require_fixture_hash
   1094             )
   1095         if scenario.leader_burn_minimum_height is not None:
   1096             assert_leader_uses_burn_from_height(
   1097                 scenario.through, scenario.leader_burn_minimum_height
   1098             )
   1099         assert_api_health(statuses, scenario.through)
   1100         print(f"e2e scenario {name} passed")
   1101     except Exception:
   1102         compose("logs", "--tail", "200", *SERVICES, check=False)
   1103         raise
   1104     finally:
   1105         if not keep:
   1106             compose("down", "--remove-orphans", check=False)
   1107 
   1108 
   1109 def run_sync_resilience_scenario(
   1110     timeout: float, build: bool, keep: bool, evidence_dir: Path | None
   1111 ) -> None:
   1112     name = "sync-resilience"
   1113     service = SYNC_SERVICE
   1114     sync_services = (*SERVICES, service)
   1115     evidence_run, evidence = create_evidence_run(evidence_dir, name)
   1116     print(
   1117         "running e2e scenario sync-resilience: interrupted empty bootstrap and stale range sync",
   1118         flush=True,
   1119     )
   1120     previous_reference_mining = os.environ.get(REFERENCE_MINING_ENV)
   1121     os.environ[REFERENCE_MINING_ENV] = "false"
   1122     try:
   1123         # Sync interruption is the variable under test. Start the reference chain
   1124         # paused so a slow deployment host cannot turn range validation into an
   1125         # unrelated moving-tip fork race while the seventh node catches up.
   1126         restore_snapshot("first-objective-checkpoint")
   1127         start(build)
   1128         initial = wait_for_height(1_001, timeout, converge=True)
   1129         initial_height = min(
   1130             status["chain"]["height"] for status in initial.values()
   1131         )
   1132         evidence["phases"]["initial"] = evidence_statuses(initial)
   1133         write_evidence_report(evidence_run, evidence)
   1134 
   1135         compose("rm", "--force", "--stop", service, check=False)
   1136         _OPENERS.pop(service, None)
   1137         sync_directory = runtime_dir() / service
   1138         if sync_directory.exists():
   1139             shutil.rmtree(sync_directory)
   1140         sync_directory.mkdir(parents=True)
   1141         chain_database = sync_directory / "chain.sqlite3"
   1142         evidence["phases"]["empty_bootstrap_started"] = {
   1143             "service": service,
   1144             "source_height": initial_height,
   1145             "chain_snapshot_present": False,
   1146         }
   1147         write_evidence_report(evidence_run, evidence)
   1148         compose("up", "--detach", "--no-deps", service)
   1149         time.sleep(0.05)
   1150         compose("kill", "--signal", "SIGKILL", service)
   1151         _OPENERS.pop(service, None)
   1152         try:
   1153             interrupted = read_chain_metadata(chain_database)
   1154         except (E2EError, sqlite3.Error):
   1155             interrupted = None
   1156         if interrupted is not None and interrupted["height"] >= initial_height:
   1157             raise E2EError(
   1158                 "empty bootstrap completed before it could be interrupted; "
   1159                 "increase the fixture height"
   1160             )
   1161         evidence["phases"]["empty_bootstrap_interrupted"] = {
   1162             "signal": "SIGKILL",
   1163             "persisted_height": (
   1164                 interrupted["height"] if interrupted is not None else None
   1165             ),
   1166         }
   1167         write_evidence_report(evidence_run, evidence)
   1168         compose("start", service)
   1169         empty_synced = wait_for_height(
   1170             initial_height, timeout, converge=True, services=sync_services
   1171         )
   1172         assert_converged(empty_synced)
   1173         evidence["phases"]["empty_bootstrap_resumed"] = {
   1174             "interrupted_before_target_persisted": True,
   1175             "nodes": evidence_statuses(empty_synced),
   1176         }
   1177         write_evidence_report(evidence_run, evidence)
   1178 
   1179         stale_target = min(empty_synced[node]["chain"]["height"] for node in SERVICES)
   1180         compose("stop", service)
   1181         _OPENERS.pop(service, None)
   1182         stale = restore_node_databases(
   1183             service, "pre-fallback-invalidation", fixture_service="node6"
   1184         )
   1185         if stale["height"] >= stale_target:
   1186             raise E2EError(
   1187                 f"stale fixture height {stale['height']} is not below target {stale_target}"
   1188             )
   1189         evidence["phases"]["stale_range_started"] = {
   1190             "service": service,
   1191             "stale_snapshot": "pre-fallback-invalidation",
   1192             "stale_height": stale["height"],
   1193             "stale_tip_hash": stale["tip_hash"],
   1194             "minimum_target_height": stale_target,
   1195         }
   1196         write_evidence_report(evidence_run, evidence)
   1197         compose("up", "--detach", "--no-deps", service)
   1198         _OPENERS.pop(service, None)
   1199         progress = wait_for_active_sync(service, stale_target, timeout)
   1200         compose("kill", "--signal", "SIGKILL", service)
   1201         _OPENERS.pop(service, None)
   1202         interrupted = read_chain_metadata(chain_database)
   1203         if interrupted["height"] >= progress["target_height"]:
   1204             raise E2EError(
   1205                 "stale range sync reached its target before process interruption"
   1206             )
   1207         evidence["phases"]["stale_range_interrupted"] = {
   1208             **progress,
   1209             "signal": "SIGKILL",
   1210             "persisted_height": interrupted["height"],
   1211             "persisted_tip_hash": interrupted["tip_hash"],
   1212         }
   1213         write_evidence_report(evidence_run, evidence)
   1214         compose("start", service)
   1215         stale_synced = wait_for_height(
   1216             stale_target, timeout, converge=True, services=sync_services
   1217         )
   1218         assert_converged(stale_synced)
   1219         assert_api_health(stale_synced, stale_target, services=sync_services)
   1220         evidence["phases"]["stale_range_resumed"] = {
   1221             "nodes": evidence_statuses(stale_synced),
   1222         }
   1223         evidence["outcome"] = "passed"
   1224         evidence["finished_at"] = datetime.now(timezone.utc).isoformat()
   1225         write_evidence_report(evidence_run, evidence)
   1226         capture_evidence_logs(evidence_run, sync_services)
   1227         print(
   1228             f"e2e scenario {name} passed: empty bootstrap and range sync "
   1229             f"resumed through at least height {stale_target}",
   1230             flush=True,
   1231         )
   1232     except Exception as error:
   1233         evidence["outcome"] = "failed"
   1234         evidence["finished_at"] = datetime.now(timezone.utc).isoformat()
   1235         evidence["error"] = {
   1236             "type": type(error).__name__,
   1237             "message": str(error),
   1238         }
   1239         write_evidence_report(evidence_run, evidence)
   1240         capture_evidence_logs(evidence_run, sync_services)
   1241         compose("logs", "--tail", "300", *sync_services, check=False)
   1242         raise
   1243     finally:
   1244         try:
   1245             if not keep:
   1246                 compose("down", "--remove-orphans", check=False)
   1247         finally:
   1248             if previous_reference_mining is None:
   1249                 os.environ.pop(REFERENCE_MINING_ENV, None)
   1250             else:
   1251                 os.environ[REFERENCE_MINING_ENV] = previous_reference_mining
   1252 
   1253 
   1254 def run_partition_recovery_scenario(
   1255     timeout: float, build: bool, keep: bool, evidence_dir: Path | None
   1256 ) -> None:
   1257     name = "partition-recovery"
   1258     evidence_run, evidence = create_evidence_run(evidence_dir, name)
   1259     print(
   1260         "running e2e scenario partition-recovery: physical 3-3 split, heal and restart",
   1261         flush=True,
   1262     )
   1263     mining_settings = None
   1264     restore_snapshot("first-objective-checkpoint")
   1265     try:
   1266         start(build)
   1267         initial = wait_for_height(1_001, timeout, converge=True)
   1268         evidence["phases"]["initial"] = evidence_statuses(initial)
   1269         write_evidence_report(evidence_run, evidence)
   1270         mining_settings = automatic_finalization_settings(initial)
   1271         recovery_workers = configure_partition_recovery_workers(
   1272             initial, mining_settings, timeout
   1273         )
   1274         apply_partition()
   1275         partitioned_start = all_statuses()
   1276         left, right = PARTITION_GROUPS
   1277         partition_boundaries = {
   1278             label: max(
   1279                 partitioned_start[service]["chain"]["height"] for service in group
   1280             )
   1281             for label, group in (("left", left), ("right", right))
   1282         }
   1283         evidence["phases"]["partition_started"] = {
   1284             "boundaries": partition_boundaries,
   1285             "recovery_workers": recovery_workers,
   1286             "nodes": evidence_statuses(partitioned_start),
   1287         }
   1288         write_evidence_report(evidence_run, evidence)
   1289         partitioned, recovery_heights = wait_for_partition_recovery(
   1290             partition_boundaries, timeout
   1291         )
   1292         evidence["phases"]["partition_recovery"] = {
   1293             "recovery_heights": recovery_heights,
   1294             "nodes": evidence_statuses(partitioned),
   1295         }
   1296         write_evidence_report(evidence_run, evidence)
   1297         print(
   1298             "partition recovery observed: "
   1299             + json.dumps(recovery_heights, sort_keys=True),
   1300             flush=True,
   1301         )
   1302 
   1303         right_worker = recovery_workers["right"]
   1304         right_mining = mining_settings[right_worker]
   1305         node_form(
   1306             right_worker,
   1307             "/api/settings/burn-per-block",
   1308             {
   1309                 "enabled": "false",
   1310                 "amount": right_mining["amount"],
   1311                 "fee_per_byte": right_mining["fee_per_byte"],
   1312             },
   1313         )
   1314         # Stop the competing branch before reconnecting the islands. Otherwise
   1315         # equally paced recovery workers can keep both branches growing forever.
   1316         time.sleep(2)
   1317         clear_partition()
   1318         heal_target = max(status["chain"]["height"] for status in partitioned.values())
   1319         healed = wait_for_height(heal_target, timeout, converge=True)
   1320         assert_converged(healed)
   1321         restore_automatic_finalization(mining_settings)
   1322         mining_settings = None
   1323         partition_height = max(partition_boundaries.values())
   1324         canonical_recoveries = [
   1325             block
   1326             for block in recent_blocks(SERVICES[0])
   1327             if block.get("finalizer_mode") == "recovery"
   1328             and int(block.get("height", -1)) > partition_height
   1329         ]
   1330         if not canonical_recoveries:
   1331             raise E2EError(
   1332                 "healed canonical chain contains no partition recovery block"
   1333             )
   1334         canonical_recovery_height = max(
   1335             int(block["height"]) for block in canonical_recoveries
   1336         )
   1337         canonical_recovery = next(
   1338             block
   1339             for block in canonical_recoveries
   1340             if int(block["height"]) == canonical_recovery_height
   1341         )
   1342         evidence["phases"]["healed"] = {
   1343             "canonical_recovery": evidence_block(canonical_recovery),
   1344             "nodes": evidence_statuses(healed),
   1345         }
   1346         write_evidence_report(evidence_run, evidence)
   1347 
   1348         restart_height = max(status["chain"]["height"] for status in healed.values())
   1349         compose("restart", "node6")
   1350         _OPENERS.pop("node6", None)
   1351         evidence["phases"]["restart"] = {
   1352             "service": "node6",
   1353             "after_height": restart_height,
   1354         }
   1355         write_evidence_report(evidence_run, evidence)
   1356         resumed, ticket = wait_for_ticket_after(
   1357             max(canonical_recovery_height, restart_height), timeout
   1358         )
   1359         assert_converged(resumed)
   1360         assert_api_health(resumed, int(ticket["height"]))
   1361         evidence["phases"]["resumed"] = {
   1362             "ticket": evidence_block(ticket),
   1363             "nodes": evidence_statuses(resumed),
   1364         }
   1365         evidence["outcome"] = "passed"
   1366         evidence["finished_at"] = datetime.now(timezone.utc).isoformat()
   1367         write_evidence_report(evidence_run, evidence)
   1368         capture_evidence_logs(evidence_run)
   1369         print(
   1370             f"e2e scenario {name} passed: recovery at {canonical_recovery_height}, "
   1371             f"rank-0 ticket resumed at {ticket['height']}",
   1372             flush=True,
   1373         )
   1374     except Exception as error:
   1375         evidence["outcome"] = "failed"
   1376         evidence["finished_at"] = datetime.now(timezone.utc).isoformat()
   1377         evidence["error"] = {
   1378             "type": type(error).__name__,
   1379             "message": str(error),
   1380         }
   1381         write_evidence_report(evidence_run, evidence)
   1382         capture_evidence_logs(evidence_run)
   1383         compose("logs", "--tail", "300", *SERVICES, check=False)
   1384         raise
   1385     finally:
   1386         clear_partition()
   1387         if mining_settings is not None:
   1388             try:
   1389                 restore_automatic_finalization(mining_settings)
   1390             except Exception as error:
   1391                 print(
   1392                     f"warning: could not restore automatic finalization settings: {error}",
   1393                     file=sys.stderr,
   1394                 )
   1395         if not keep:
   1396             compose("down", "--remove-orphans", check=False)
   1397 
   1398 
   1399 def run_tests(
   1400     name: str, timeout: float, build: bool, keep: bool, evidence_dir: Path | None
   1401 ) -> None:
   1402     if name in ("snapshots", "all"):
   1403         test_snapshots()
   1404     if name == "all":
   1405         selected = [*SCENARIOS, *SPECIAL_SCENARIOS]
   1406     elif name == "post-activation":
   1407         test_snapshots()
   1408         selected = list(POST_ACTIVATION_SCENARIOS)
   1409     else:
   1410         selected = [name]
   1411     selected = [scenario_name for scenario_name in selected if scenario_name != "snapshots"]
   1412     for index, scenario_name in enumerate(selected):
   1413         scenario_keep = keep and index == len(selected) - 1
   1414         if scenario_name == "partition-recovery":
   1415             run_partition_recovery_scenario(
   1416                 timeout,
   1417                 build and index == 0,
   1418                 scenario_keep,
   1419                 evidence_dir,
   1420             )
   1421         elif scenario_name == "sync-resilience":
   1422             run_sync_resilience_scenario(
   1423                 timeout,
   1424                 build and index == 0,
   1425                 scenario_keep,
   1426                 evidence_dir,
   1427             )
   1428         else:
   1429             run_scenario(
   1430                 scenario_name,
   1431                 SCENARIOS[scenario_name],
   1432                 timeout,
   1433                 build and index == 0,
   1434                 scenario_keep,
   1435             )
   1436     print(f"e2e test run passed: {name}")
   1437 
   1438 
   1439 def parser() -> argparse.ArgumentParser:
   1440     result = argparse.ArgumentParser(description=__doc__)
   1441     commands = result.add_subparsers(dest="command", required=True)
   1442 
   1443     up = commands.add_parser("up", help="start the isolated six-node e2e network")
   1444     up.add_argument("--build", action="store_true")
   1445     up.add_argument("--snapshot", help="restore this checkpoint before starting")
   1446 
   1447     commands.add_parser("down", help="stop the e2e network without deleting runtime data")
   1448     commands.add_parser("reset", help="stop the network and delete mutable e2e node data")
   1449     commands.add_parser("status", help="show authenticated live status for every node")
   1450 
   1451     wait = commands.add_parser("wait", help="wait until every node reaches a height")
   1452     wait.add_argument("height", type=int)
   1453     wait.add_argument("--timeout", type=float, default=600)
   1454     wait.add_argument("--converge", action="store_true")
   1455 
   1456     capture = commands.add_parser("capture", help="capture all node identities and chain DBs")
   1457     capture.add_argument("name")
   1458     capture.add_argument("--height", type=int, required=True)
   1459     capture.add_argument("--timeout", type=float, default=7_200)
   1460     capture.add_argument("--force", action="store_true")
   1461 
   1462     restore = commands.add_parser("restore", help="replace runtime data with a checkpoint")
   1463     restore.add_argument("name")
   1464 
   1465     plan = commands.add_parser("capture-plan", help="capture each checkpoint in a JSON plan")
   1466     plan.add_argument("--plan", type=Path, default=E2E_DIR / "checkpoints.json")
   1467     plan.add_argument("--timeout", type=float, default=7_200)
   1468     plan.add_argument("--force", action="store_true")
   1469 
   1470     smoke_test = commands.add_parser(
   1471         "smoke", help="restore a checkpoint and test network convergence"
   1472     )
   1473     smoke_test.add_argument("--from", dest="snapshot", required=True)
   1474     smoke_test.add_argument("--through", type=int, required=True)
   1475     smoke_test.add_argument("--timeout", type=float, default=600)
   1476     smoke_test.add_argument("--build", action="store_true")
   1477     smoke_test.add_argument("--keep", action="store_true")
   1478 
   1479     tests = commands.add_parser("test", help="run named e2e assertions")
   1480     tests.add_argument(
   1481         "scenario",
   1482         nargs="?",
   1483         default="all",
   1484         choices=(
   1485             "all",
   1486             "post-activation",
   1487             "snapshots",
   1488             *SCENARIOS,
   1489             *SPECIAL_SCENARIOS,
   1490         ),
   1491     )
   1492     tests.add_argument("--timeout", type=float, default=600)
   1493     tests.add_argument("--build", action="store_true")
   1494     tests.add_argument("--keep", action="store_true")
   1495     tests.add_argument(
   1496         "--evidence-dir",
   1497         type=Path,
   1498         help="write scenario phase reports and node logs below this directory",
   1499     )
   1500     return result
   1501 
   1502 
   1503 def main() -> int:
   1504     args = parser().parse_args()
   1505     try:
   1506         if args.command == "up":
   1507             if args.snapshot:
   1508                 restore_snapshot(args.snapshot)
   1509             start(args.build)
   1510         elif args.command == "down":
   1511             compose("down", "--remove-orphans")
   1512         elif args.command == "reset":
   1513             reset_runtime()
   1514         elif args.command == "status":
   1515             statuses = all_statuses()
   1516             assert_e2e_profile(statuses)
   1517             print(json.dumps({name: compact_status(value) for name, value in statuses.items()}, indent=2))
   1518         elif args.command == "wait":
   1519             statuses = wait_for_height(args.height, args.timeout, args.converge)
   1520             print(json.dumps({name: compact_status(value) for name, value in statuses.items()}, indent=2))
   1521         elif args.command == "capture":
   1522             capture_snapshot(args.name, args.height, args.timeout, args.force)
   1523         elif args.command == "restore":
   1524             restore_snapshot(args.name)
   1525         elif args.command == "capture-plan":
   1526             capture_plan(args.plan, args.timeout, args.force)
   1527         elif args.command == "smoke":
   1528             smoke(args.snapshot, args.through, args.timeout, args.build, args.keep)
   1529         elif args.command == "test":
   1530             run_tests(
   1531                 args.scenario,
   1532                 args.timeout,
   1533                 args.build,
   1534                 args.keep,
   1535                 args.evidence_dir,
   1536             )
   1537         return 0
   1538     except (E2EError, OSError, sqlite3.Error, subprocess.CalledProcessError) as error:
   1539         print(f"e2e error: {error}", file=sys.stderr)
   1540         return 1
   1541 
   1542 
   1543 if __name__ == "__main__":
   1544     raise SystemExit(main())