From caf1663158bec193ee5d059486bccc2d1a6800b4 Mon Sep 17 00:00:00 2001 From: Yostra Date: Fri, 25 Dec 2020 20:06:12 -0500 Subject: [PATCH] fetch recent block don't enter sync mode --- src/simulator/full_node_simulator.py | 1 - src/wallet/wallet_node.py | 99 +++++++++++++++++++++++----- 2 files changed, 82 insertions(+), 18 deletions(-) diff --git a/src/simulator/full_node_simulator.py b/src/simulator/full_node_simulator.py index a488b18339..f2ae7964fd 100644 --- a/src/simulator/full_node_simulator.py +++ b/src/simulator/full_node_simulator.py @@ -65,7 +65,6 @@ class FullNodeSimulator(FullNodeAPI): farmer_reward_puzzle_hash=target, pool_reward_puzzle_hash=target, block_list_input=current_blocks, - force_overflow=True, guarantee_block=True, ) rr = RespondSubBlock(more[-1]) diff --git a/src/wallet/wallet_node.py b/src/wallet/wallet_node.py index 877dece42a..790ba9934f 100644 --- a/src/wallet/wallet_node.py +++ b/src/wallet/wallet_node.py @@ -104,6 +104,7 @@ class WalletNode: self.server = None self.wsm_close_task = None self.sync_task: Optional[Task] = None + self.new_peak_lock = asyncio.Lock() def get_key_for_fingerprint(self, fingerprint): private_keys = self.keychain.get_all_private_keys() @@ -323,24 +324,90 @@ class WalletNode: return True return False + async def complete_blocks(self, header_blocks: List[HeaderBlock], peer: WSChiaConnection): + header_block_records: List[HeaderBlockRecord] = [] + for block in header_blocks: + if block.is_block: + # Find additions and removals + ( + additions, + removals, + ) = await self.wallet_state_manager.get_filter_additions_removals(block, block.transactions_filter) + + # Get Additions + added_coins = await self.get_additions(peer, block, additions) + if added_coins is None: + raise ValueError("Failed to fetch additions") + + # Get removals + removed_coins = await self.get_removals(peer, block, added_coins, removals) + if removed_coins is None: + raise ValueError("Failed to fetch removals") + hbr = HeaderBlockRecord(block, added_coins, removed_coins) + else: + hbr = HeaderBlockRecord(block, [], []) + header_block_records.append(hbr) + ( + result, + error, + fork_h, + ) = await self.wallet_state_manager.blockchain.receive_block(hbr) + if result == ReceiveBlockResult.NEW_PEAK: + self.wallet_state_manager.state_changed("new_block") + elif result == ReceiveBlockResult.INVALID_BLOCK: + self.log.info(f"Invalid block from peer: {peer.get_peer_info()}") + await peer.close() + return + async def new_peak(self, peak: wallet_protocol.NewPeak, peer: WSChiaConnection): if self.wallet_state_manager is None: return curr_peak = self.wallet_state_manager.blockchain.get_peak() - if curr_peak is not None and curr_peak.weight > curr_peak.weight: + if curr_peak is not None and curr_peak.weight >= peak.weight: return + async with self.new_peak_lock: + request = wallet_protocol.RequestSubBlockHeader(peak.sub_block_height) + response: Optional[RespondSubBlockHeader] = await peer.request_sub_block_header(request) - request = wallet_protocol.RequestSubBlockHeader(peak.sub_block_height) - response: Optional[RespondSubBlockHeader] = await peer.request_sub_block_header(request) + if ( + response is not None + and isinstance(response, RespondSubBlockHeader) + and response.header_block is not None + ): + header_block = response.header_block - if response is not None and response.header_block is not None: - hb = response.header_block - - if curr_peak is not None and hb.prev_header_hash != peak.header_hash: - # only request weight proofs if we are past the first sub epoch - if hb.sub_block_height > self.constants.SUB_EPOCH_SUB_BLOCKS: - weight_request = RequestProofOfWeight(hb.sub_block_height, hb.header_hash) + if ( + curr_peak is None and header_block.sub_block_height < self.constants.WEIGHT_PROOF_RECENT_BLOCKS + ) or ( + curr_peak is not None + and curr_peak.height > header_block.sub_block_height - self.constants.WEIGHT_PROOF_RECENT_BLOCKS + ): + top = header_block + blocks = [top] + # Fetch blocks backwards until we hit the one that we have, + # then complete them with additions / removals going forward + while ( + top.prev_header_hash not in self.wallet_state_manager.blockchain.sub_blocks + and top.sub_block_height > 0 + ): + request_prev = wallet_protocol.RequestSubBlockHeader(top.sub_block_height - 1) + response_prev: Optional[RespondSubBlockHeader] = await peer.request_sub_block_header( + request_prev + ) + if response_prev is None: + return + if not isinstance(response_prev, RespondSubBlockHeader): + return + prev_head = response_prev.header_block + blocks.append(prev_head) + top = prev_head + blocks.reverse() + await self.complete_blocks(blocks, peer) + else: + # Request weight proof + # Sync if PoW validates + weight_request = RequestProofOfWeight(header_block.sub_block_height, header_block.header_hash) weight_proof_response: RespondProofOfWeight = await peer.request_proof_of_weight(weight_request) if weight_proof_response is None: return @@ -355,9 +422,11 @@ class WalletNode: ) return None self.log.info(f"Validated, fork point is {fork_point}") - self.wallet_state_manager.sync_store.add_potential_fork_point(hb.header_hash, uint32(fork_point)) - self.wallet_state_manager.sync_store.add_potential_peak(hb) - self.start_sync() + self.wallet_state_manager.sync_store.add_potential_fork_point( + header_block.header_hash, uint32(fork_point) + ) + self.wallet_state_manager.sync_store.add_potential_peak(header_block) + self.start_sync() def start_sync(self): self.log.info("self.sync_event.set()") @@ -430,8 +499,6 @@ class WalletNode: self.log.info("No peers to sync to") return - fetched_blocks: Dict[int, HeaderBlockRecord] = {} - fork_height = self.wallet_state_manager.sync_store.get_potential_fork_point(peak.header_hash) if fork_height is None: fork_height = 0 @@ -469,8 +536,6 @@ class WalletNode: else: header_block_record = HeaderBlockRecord(block_i, [], []) - fetched_blocks[i] = header_block_record - self.log.info(f"Adding {header_block_record.additions}") ( result, error,