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())