mirror of
https://github.com/Chia-Network/chia-blockchain.git
synced 2026-09-05 10:05:00 -05:00
* protocols|server: Define new harvester plot refreshing protocol messages * protocols: Bump `protocol_version` to `0.0.34` * tests: Introduce `setup_farmer_multi_harvester` Allows to run a test setup with 1 farmer and mutiple harvesters. * plotting: Add an initial plot loading indication to `PlotManager` * plotting|tests: Don't add removed duplicates to `total_result.removed` `PlotRefreshResult.removed` should only contain plots that were loaded properly before they were removed. It shouldn't contain e.g. removed duplicates or invalid plots since those are synced in an extra sync step and not as diff but as whole list every time. * harvester: Reset `PlotManager` on shutdown * plot_sync: Implement plot sync protocol * farmer|harvester: Integrate and enable plot sync * tests: Implement tests for the plot sync protocol * farmer|tests: Drop obsolete harvester caching code * setup: Add `chia.plot_sync` to packages * plot_sync: Type hints in `DeltaType` * plot_sync: Drop parameters in `super()` calls * plot_sync: Introduce `send_response` helper in `Receiver._process` * plot_sync: Add some parentheses Co-authored-by: Kyle Altendorf <sda@fstab.net> * plot_sync: Additional hint for a `Receiver.process_path_list` parameter * plot_sync: Force named parameters in `Receiver.process_path_list` * test: Fix fixtures after rebase * tests: Fix sorting after rebase * tests: Return type hint for `plot_sync_setup` * tests: Rename `WSChiaConnection` and move it in the outer scope * tests|plot_sync: More type hints * tests: Rework some delta tests * tests: Drop a `range` and iterate over the list directly * tests: Use the proper flags to overwrite * test: More missing duplicates tests * tests: Drop `ExpectedResult.reset` * tests: Reduce some asserts * tests: Add messages to some `assert False` statements * tests: Introduce `ErrorSimulation` enum in `test_sync_simulated.py` * tests: Use `secrects` instead of `Crypto.Random` * Fixes after rebase * Import from `typing_extensions` to support python 3.7 * Drop task name to support python 3.7 * Introduce `Sender.syncing`, `Sender.connected` and a log about the task * Add `tests/plot_sync/config.py` * Align the multi harvester fixture with what we do in other places * Update the workflows Co-authored-by: Kyle Altendorf <sda@fstab.net>
342 lines
11 KiB
Python
342 lines
11 KiB
Python
import asyncio
|
|
import logging
|
|
import signal
|
|
import sqlite3
|
|
from pathlib import Path
|
|
from secrets import token_bytes
|
|
from typing import AsyncGenerator, Optional
|
|
|
|
from chia.consensus.constants import ConsensusConstants
|
|
from chia.daemon.server import WebSocketServer, daemon_launch_lock_path, singleton
|
|
from chia.server.start_farmer import service_kwargs_for_farmer
|
|
from chia.server.start_full_node import service_kwargs_for_full_node
|
|
from chia.server.start_harvester import service_kwargs_for_harvester
|
|
from chia.server.start_introducer import service_kwargs_for_introducer
|
|
from chia.server.start_service import Service
|
|
from chia.server.start_timelord import service_kwargs_for_timelord
|
|
from chia.server.start_wallet import service_kwargs_for_wallet
|
|
from chia.simulator.start_simulator import service_kwargs_for_full_node_simulator
|
|
from chia.timelord.timelord_launcher import kill_processes, spawn_process
|
|
from chia.util.bech32m import encode_puzzle_hash
|
|
from chia.util.config import load_config, save_config
|
|
from chia.util.ints import uint16
|
|
from chia.util.keychain import bytes_to_mnemonic
|
|
from tests.block_tools import BlockTools
|
|
from tests.util.keyring import TempKeyring
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
|
|
async def setup_daemon(btools: BlockTools) -> AsyncGenerator[WebSocketServer, None]:
|
|
root_path = btools.root_path
|
|
config = btools.config
|
|
assert "daemon_port" in config
|
|
lockfile = singleton(daemon_launch_lock_path(root_path))
|
|
crt_path = root_path / config["daemon_ssl"]["private_crt"]
|
|
key_path = root_path / config["daemon_ssl"]["private_key"]
|
|
ca_crt_path = root_path / config["private_ssl_ca"]["crt"]
|
|
ca_key_path = root_path / config["private_ssl_ca"]["key"]
|
|
assert lockfile is not None
|
|
shutdown_event = asyncio.Event()
|
|
ws_server = WebSocketServer(root_path, ca_crt_path, ca_key_path, crt_path, key_path, shutdown_event)
|
|
await ws_server.start()
|
|
|
|
yield ws_server
|
|
|
|
await ws_server.stop()
|
|
|
|
|
|
async def setup_full_node(
|
|
consensus_constants: ConsensusConstants,
|
|
db_name,
|
|
self_hostname: str,
|
|
port,
|
|
rpc_port,
|
|
local_bt: BlockTools,
|
|
introducer_port=None,
|
|
simulator=False,
|
|
send_uncompact_interval=0,
|
|
sanitize_weight_proof_only=False,
|
|
connect_to_daemon=False,
|
|
db_version=1,
|
|
):
|
|
db_path = local_bt.root_path / f"{db_name}"
|
|
if db_path.exists():
|
|
db_path.unlink()
|
|
|
|
if db_version > 1:
|
|
with sqlite3.connect(db_path) as connection:
|
|
connection.execute("CREATE TABLE database_version(version int)")
|
|
connection.execute("INSERT INTO database_version VALUES (?)", (db_version,))
|
|
connection.commit()
|
|
|
|
if connect_to_daemon:
|
|
assert local_bt.config["daemon_port"] is not None
|
|
config = local_bt.config["full_node"]
|
|
|
|
config["database_path"] = db_name
|
|
config["send_uncompact_interval"] = send_uncompact_interval
|
|
config["target_uncompact_proofs"] = 30
|
|
config["peer_connect_interval"] = 50
|
|
config["sanitize_weight_proof_only"] = sanitize_weight_proof_only
|
|
if introducer_port is not None:
|
|
config["introducer_peer"]["host"] = self_hostname
|
|
config["introducer_peer"]["port"] = introducer_port
|
|
else:
|
|
config["introducer_peer"] = None
|
|
config["dns_servers"] = []
|
|
config["port"] = port
|
|
config["rpc_port"] = rpc_port
|
|
overrides = config["network_overrides"]["constants"][config["selected_network"]]
|
|
updated_constants = consensus_constants.replace_str_to_bytes(**overrides)
|
|
if simulator:
|
|
kwargs = service_kwargs_for_full_node_simulator(local_bt.root_path, config, local_bt)
|
|
else:
|
|
kwargs = service_kwargs_for_full_node(local_bt.root_path, config, updated_constants)
|
|
|
|
kwargs.update(
|
|
parse_cli_args=False,
|
|
connect_to_daemon=connect_to_daemon,
|
|
service_name_prefix="test_",
|
|
)
|
|
|
|
service = Service(**kwargs, running_new_process=False)
|
|
|
|
await service.start()
|
|
|
|
yield service._api
|
|
|
|
service.stop()
|
|
await service.wait_closed()
|
|
if db_path.exists():
|
|
db_path.unlink()
|
|
|
|
|
|
# Note: convert these setup functions to fixtures, or push it one layer up,
|
|
# keeping these usable independently?
|
|
async def setup_wallet_node(
|
|
self_hostname: str,
|
|
port,
|
|
rpc_port,
|
|
consensus_constants: ConsensusConstants,
|
|
local_bt: BlockTools,
|
|
full_node_port=None,
|
|
introducer_port=None,
|
|
key_seed=None,
|
|
starting_height=None,
|
|
initial_num_public_keys=5,
|
|
):
|
|
with TempKeyring(populate=True) as keychain:
|
|
config = local_bt.config["wallet"]
|
|
config["port"] = port
|
|
config["rpc_port"] = rpc_port
|
|
if starting_height is not None:
|
|
config["starting_height"] = starting_height
|
|
config["initial_num_public_keys"] = initial_num_public_keys
|
|
|
|
entropy = token_bytes(32)
|
|
if key_seed is None:
|
|
key_seed = entropy
|
|
keychain.add_private_key(bytes_to_mnemonic(key_seed), "")
|
|
first_pk = keychain.get_first_public_key()
|
|
assert first_pk is not None
|
|
db_path_key_suffix = str(first_pk.get_fingerprint())
|
|
db_name = f"test-wallet-db-{port}-KEY.sqlite"
|
|
db_path_replaced: str = db_name.replace("KEY", db_path_key_suffix)
|
|
db_path = local_bt.root_path / db_path_replaced
|
|
|
|
if db_path.exists():
|
|
db_path.unlink()
|
|
config["database_path"] = str(db_name)
|
|
config["testing"] = True
|
|
|
|
config["introducer_peer"]["host"] = self_hostname
|
|
if introducer_port is not None:
|
|
config["introducer_peer"]["port"] = introducer_port
|
|
config["peer_connect_interval"] = 10
|
|
else:
|
|
config["introducer_peer"] = None
|
|
|
|
if full_node_port is not None:
|
|
config["full_node_peer"] = {}
|
|
config["full_node_peer"]["host"] = self_hostname
|
|
config["full_node_peer"]["port"] = full_node_port
|
|
else:
|
|
del config["full_node_peer"]
|
|
|
|
kwargs = service_kwargs_for_wallet(local_bt.root_path, config, consensus_constants, keychain)
|
|
kwargs.update(
|
|
parse_cli_args=False,
|
|
connect_to_daemon=False,
|
|
service_name_prefix="test_",
|
|
)
|
|
|
|
service = Service(**kwargs, running_new_process=False)
|
|
|
|
await service.start()
|
|
|
|
yield service._node, service._node.server
|
|
|
|
service.stop()
|
|
await service.wait_closed()
|
|
if db_path.exists():
|
|
db_path.unlink()
|
|
keychain.delete_all_keys()
|
|
|
|
|
|
async def setup_harvester(
|
|
root_path: Path,
|
|
self_hostname: str,
|
|
port,
|
|
rpc_port,
|
|
farmer_port,
|
|
consensus_constants: ConsensusConstants,
|
|
start_service: bool = True,
|
|
):
|
|
config = load_config(root_path, "config.yaml")
|
|
config["harvester"]["port"] = port
|
|
config["harvester"]["rpc_port"] = rpc_port
|
|
config["harvester"]["farmer_peer"]["host"] = self_hostname
|
|
config["harvester"]["farmer_peer"]["port"] = farmer_port
|
|
save_config(root_path, "config.yaml", config)
|
|
kwargs = service_kwargs_for_harvester(root_path, config["harvester"], consensus_constants)
|
|
kwargs.update(
|
|
parse_cli_args=False,
|
|
connect_to_daemon=False,
|
|
service_name_prefix="test_",
|
|
)
|
|
|
|
service = Service(**kwargs, running_new_process=False)
|
|
|
|
if start_service:
|
|
await service.start()
|
|
|
|
yield service
|
|
|
|
service.stop()
|
|
await service.wait_closed()
|
|
|
|
|
|
async def setup_farmer(
|
|
b_tools: BlockTools,
|
|
self_hostname: str,
|
|
port,
|
|
rpc_port,
|
|
consensus_constants: ConsensusConstants,
|
|
full_node_port: Optional[uint16] = None,
|
|
start_service: bool = True,
|
|
):
|
|
config = b_tools.config["farmer"]
|
|
config_pool = b_tools.config["pool"]
|
|
|
|
config["xch_target_address"] = encode_puzzle_hash(b_tools.farmer_ph, "xch")
|
|
config["pool_public_keys"] = [bytes(pk).hex() for pk in b_tools.pool_pubkeys]
|
|
config["port"] = port
|
|
config["rpc_port"] = rpc_port
|
|
config_pool["xch_target_address"] = encode_puzzle_hash(b_tools.pool_ph, "xch")
|
|
|
|
if full_node_port:
|
|
config["full_node_peer"]["host"] = self_hostname
|
|
config["full_node_peer"]["port"] = full_node_port
|
|
else:
|
|
del config["full_node_peer"]
|
|
|
|
kwargs = service_kwargs_for_farmer(
|
|
b_tools.root_path, config, config_pool, consensus_constants, b_tools.local_keychain
|
|
)
|
|
kwargs.update(
|
|
parse_cli_args=False,
|
|
connect_to_daemon=False,
|
|
service_name_prefix="test_",
|
|
)
|
|
|
|
service = Service(**kwargs, running_new_process=False)
|
|
|
|
if start_service:
|
|
await service.start()
|
|
|
|
yield service
|
|
|
|
service.stop()
|
|
await service.wait_closed()
|
|
|
|
|
|
async def setup_introducer(bt: BlockTools, port):
|
|
kwargs = service_kwargs_for_introducer(
|
|
bt.root_path,
|
|
bt.config["introducer"],
|
|
)
|
|
kwargs.update(
|
|
advertised_port=port,
|
|
parse_cli_args=False,
|
|
connect_to_daemon=False,
|
|
service_name_prefix="test_",
|
|
)
|
|
|
|
service = Service(**kwargs, running_new_process=False)
|
|
|
|
await service.start()
|
|
|
|
yield service._api, service._node.server
|
|
|
|
service.stop()
|
|
await service.wait_closed()
|
|
|
|
|
|
async def setup_vdf_client(bt: BlockTools, self_hostname: str, port):
|
|
vdf_task_1 = asyncio.create_task(spawn_process(self_hostname, port, 1, bt.config.get("prefer_ipv6")))
|
|
|
|
def stop():
|
|
asyncio.create_task(kill_processes())
|
|
|
|
asyncio.get_running_loop().add_signal_handler(signal.SIGTERM, stop)
|
|
asyncio.get_running_loop().add_signal_handler(signal.SIGINT, stop)
|
|
|
|
yield vdf_task_1
|
|
await kill_processes()
|
|
|
|
|
|
async def setup_vdf_clients(bt: BlockTools, self_hostname: str, port):
|
|
vdf_task_1 = asyncio.create_task(spawn_process(self_hostname, port, 1, bt.config.get("prefer_ipv6")))
|
|
vdf_task_2 = asyncio.create_task(spawn_process(self_hostname, port, 2, bt.config.get("prefer_ipv6")))
|
|
vdf_task_3 = asyncio.create_task(spawn_process(self_hostname, port, 3, bt.config.get("prefer_ipv6")))
|
|
|
|
def stop():
|
|
asyncio.create_task(kill_processes())
|
|
|
|
asyncio.get_running_loop().add_signal_handler(signal.SIGTERM, stop)
|
|
asyncio.get_running_loop().add_signal_handler(signal.SIGINT, stop)
|
|
|
|
yield vdf_task_1, vdf_task_2, vdf_task_3
|
|
|
|
await kill_processes()
|
|
|
|
|
|
async def setup_timelord(
|
|
port, full_node_port, rpc_port, vdf_port, sanitizer, consensus_constants: ConsensusConstants, b_tools: BlockTools
|
|
):
|
|
config = b_tools.config["timelord"]
|
|
config["port"] = port
|
|
config["full_node_peer"]["port"] = full_node_port
|
|
config["bluebox_mode"] = sanitizer
|
|
config["fast_algorithm"] = False
|
|
config["vdf_server"]["port"] = vdf_port
|
|
config["start_rpc_server"] = True
|
|
config["rpc_port"] = rpc_port
|
|
|
|
kwargs = service_kwargs_for_timelord(b_tools.root_path, config, consensus_constants)
|
|
kwargs.update(
|
|
parse_cli_args=False,
|
|
connect_to_daemon=False,
|
|
service_name_prefix="test_",
|
|
)
|
|
|
|
service = Service(**kwargs, running_new_process=False)
|
|
|
|
await service.start()
|
|
|
|
yield service._api, service._node.server
|
|
|
|
service.stop()
|
|
await service.wait_closed()
|