From 8fab59975fe6e8bf0bc6c13d1cb206daf51adc66 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Thu, 2 Jan 2020 16:15:27 +0100 Subject: [PATCH 01/14] Adds basic ChainMonitor to the system ChainMonitor is the actor in charge of checking new blocks. It improves from the previous zmq_subscriber by also doing polling and, therefore, making sure that no blocks are missed. Documentation and tests are still required. Partially covers #31 --- pisa/chain_monitor.py | 92 ++++++++++++++++++++++++++++++++++++ pisa/pisad.py | 12 ++++- pisa/responder.py | 23 ++------- pisa/utils/zmq_subscriber.py | 34 ------------- pisa/watcher.py | 23 ++------- 5 files changed, 112 insertions(+), 72 deletions(-) create mode 100644 pisa/chain_monitor.py delete mode 100644 pisa/utils/zmq_subscriber.py diff --git a/pisa/chain_monitor.py b/pisa/chain_monitor.py new file mode 100644 index 0000000..74b21f1 --- /dev/null +++ b/pisa/chain_monitor.py @@ -0,0 +1,92 @@ +import zmq +import binascii +from queue import Queue +from threading import Thread, Event, Condition + +from common.logger import Logger +from pisa.conf import FEED_PROTOCOL, FEED_ADDR, FEED_PORT +from pisa.block_processor import BlockProcessor + +logger = Logger("ChainMonitor") + + +class ChainMonitor: + def __init__(self): + self.best_tip = None + self.last_tips = [] + self.terminate = False + + self.check_tip = Event() + self.lock = Condition() + + self.zmqContext = zmq.Context() + self.zmqSubSocket = self.zmqContext.socket(zmq.SUB) + self.zmqSubSocket.setsockopt(zmq.RCVHWM, 0) + self.zmqSubSocket.setsockopt_string(zmq.SUBSCRIBE, "hashblock") + self.zmqSubSocket.connect("%s://%s:%s" % (FEED_PROTOCOL, FEED_ADDR, FEED_PORT)) + + self.watcher_queue = None + self.watcher_asleep = True + self.responder_queue = None + self.responder_asleep = True + + def attach_watcher(self, queue, awake): + self.watcher_queue = queue + self.watcher_asleep = awake + + def attach_responder(self, queue, awake): + self.responder_queue = queue + self.responder_asleep = awake + + def notify_subscribers(self, block_hash): + if not self.watcher_asleep: + self.watcher_queue.put(block_hash) + + if not self.responder_asleep: + self.responder_queue.put(block_hash) + + def update_state(self, block_hash): + self.best_tip = block_hash + self.last_tips.append(block_hash) + + if len(self.last_tips) > 10: + self.last_tips.pop(0) + + def monitor_chain(self): + self.best_tip = BlockProcessor.get_best_block_hash() + Thread(target=self.monitor_chain_polling).start() + Thread(target=self.monitor_chain_zmq).start() + + def monitor_chain_polling(self): + while self.terminate: + self.check_tip.wait(timeout=60) + + # Terminate could have been set wile the thread was blocked in wait + if not self.terminate: + current_tip = BlockProcessor.get_best_block_hash() + + self.lock.acquire() + if current_tip != self.best_tip: + self.update_state(current_tip) + self.notify_subscribers(current_tip) + logger.info("New block received via polling", block_hash=current_tip) + self.lock.release() + + def monitor_chain_zmq(self): + while not self.terminate: + msg = self.zmqSubSocket.recv_multipart() + + # Terminate could have been set wile the thread was blocked in recv + if not self.terminate: + topic = msg[0] + body = msg[1] + + if topic == b"hashblock": + block_hash = binascii.hexlify(body).decode("utf-8") + + self.lock.acquire() + if block_hash != self.best_tip and block_hash not in self.last_tips: + self.update_state(block_hash) + self.notify_subscribers(block_hash) + logger.info("New block received via zmq", block_hash=block_hash) + self.lock.release() diff --git a/pisa/pisad.py b/pisa/pisad.py index 9c258ce..ed185c6 100644 --- a/pisa/pisad.py +++ b/pisa/pisad.py @@ -10,6 +10,7 @@ from pisa.builder import Builder from pisa.conf import BTC_NETWORK, PISA_SECRET_KEY from pisa.responder import Responder from pisa.db_manager import DBManager +from pisa.chain_monitor import ChainMonitor from pisa.block_processor import BlockProcessor from pisa.tools import can_connect_to_bitcoind, in_correct_network @@ -46,13 +47,18 @@ if __name__ == "__main__": try: db_manager = DBManager(DB_PATH) + # Create the chain monitor and start monitoring the chain + chain_monitor = ChainMonitor() + chain_monitor.monitor_chain() + watcher_appointments_data = db_manager.load_watcher_appointments() responder_trackers_data = db_manager.load_responder_trackers() with open(PISA_SECRET_KEY, "rb") as key_file: secret_key_der = key_file.read() - watcher = Watcher(db_manager, secret_key_der) + watcher = Watcher(db_manager, chain_monitor, secret_key_der) + chain_monitor.attach_watcher(watcher.block_queue, watcher.asleep) if len(watcher_appointments_data) == 0 and len(responder_trackers_data) == 0: logger.info("Fresh bootstrap") @@ -65,7 +71,7 @@ if __name__ == "__main__": last_block_responder = db_manager.load_last_block_hash_responder() # FIXME: 32-reorgs-offline dropped txs are not used at this point. - responder = Responder(db_manager) + responder = Responder(db_manager, chain_monitor) last_common_ancestor_responder = None missed_blocks_responder = None @@ -82,6 +88,8 @@ if __name__ == "__main__": # Build Watcher with Responder and backed up data. If the blocks of both match we don't perform the # search twice. watcher.responder = responder + chain_monitor.attach_responder(responder.block_queue, responder.asleep) + if last_block_watcher is not None: if last_block_watcher == last_block_responder: missed_blocks_watcher = missed_blocks_responder diff --git a/pisa/responder.py b/pisa/responder.py index 0fce1cb..dbfac15 100644 --- a/pisa/responder.py +++ b/pisa/responder.py @@ -6,7 +6,6 @@ from common.logger import Logger from pisa.cleaner import Cleaner from pisa.carrier import Carrier from pisa.block_processor import BlockProcessor -from pisa.utils.zmq_subscriber import ZMQSubscriber CONFIRMATIONS_BEFORE_RETRY = 6 MIN_CONFIRMATIONS = 6 @@ -135,14 +134,14 @@ class Responder: """ - def __init__(self, db_manager): + def __init__(self, db_manager, chain_monitor): self.trackers = dict() self.tx_tracker_map = dict() self.unconfirmed_txs = [] self.missed_confirmations = dict() self.asleep = True self.block_queue = Queue() - self.zmq_subscriber = None + self.chain_monitor = chain_monitor self.db_manager = db_manager @staticmethod @@ -260,19 +259,8 @@ class Responder: if self.asleep: self.asleep = False - zmq_thread = Thread(target=self.do_subscribe) - responder = Thread(target=self.do_watch) - zmq_thread.start() - responder.start() - - def do_subscribe(self): - """ - Initializes a :obj:`ZMQSubscriber ` instance to listen to new blocks - from ``bitcoind``. Block ids are received trough the ``block_queue``. - """ - - self.zmq_subscriber = ZMQSubscriber(parent="Responder") - self.zmq_subscriber.handle(self.block_queue) + self.chain_monitor.responder_asleep = False + Thread(target=self.do_watch).start() def do_watch(self): """ @@ -327,8 +315,7 @@ class Responder: # Go back to sleep if there are no more pending trackers self.asleep = True - self.zmq_subscriber.terminate = True - self.block_queue = Queue() + self.chain_monitor.responder_asleep = True logger.info("No more pending trackers, going back to sleep") diff --git a/pisa/utils/zmq_subscriber.py b/pisa/utils/zmq_subscriber.py deleted file mode 100644 index ecec9af..0000000 --- a/pisa/utils/zmq_subscriber.py +++ /dev/null @@ -1,34 +0,0 @@ -import zmq -import binascii -from common.logger import Logger -from pisa.conf import FEED_PROTOCOL, FEED_ADDR, FEED_PORT - - -# ToDo: #7-add-async-back-to-zmq -class ZMQSubscriber: - """ Adapted from https://github.com/bitcoin/bitcoin/blob/master/contrib/zmq/zmq_sub.py""" - - def __init__(self, parent): - self.zmqContext = zmq.Context() - self.zmqSubSocket = self.zmqContext.socket(zmq.SUB) - self.zmqSubSocket.setsockopt(zmq.RCVHWM, 0) - self.zmqSubSocket.setsockopt_string(zmq.SUBSCRIBE, "hashblock") - self.zmqSubSocket.connect("%s://%s:%s" % (FEED_PROTOCOL, FEED_ADDR, FEED_PORT)) - self.logger = Logger("ZMQSubscriber-{}".format(parent)) - - self.terminate = False - - def handle(self, block_queue): - while not self.terminate: - msg = self.zmqSubSocket.recv_multipart() - - # Terminate could have been set wile the thread was blocked in recv - if not self.terminate: - topic = msg[0] - body = msg[1] - - if topic == b"hashblock": - block_hash = binascii.hexlify(body).decode("utf-8") - block_queue.put(block_hash) - - self.logger.info("New block received via ZMQ", block_hash=block_hash) diff --git a/pisa/watcher.py b/pisa/watcher.py index 9d659db..d6e0729 100644 --- a/pisa/watcher.py +++ b/pisa/watcher.py @@ -9,7 +9,6 @@ from common.logger import Logger from pisa.cleaner import Cleaner from pisa.responder import Responder from pisa.block_processor import BlockProcessor -from pisa.utils.zmq_subscriber import ZMQSubscriber from pisa.conf import EXPIRY_DELTA, MAX_APPOINTMENTS logger = Logger("Watcher") @@ -58,13 +57,13 @@ class Watcher: """ - def __init__(self, db_manager, sk_der, responder=None, max_appointments=MAX_APPOINTMENTS): + def __init__(self, db_manager, chain_monitor, sk_der, responder=None, max_appointments=MAX_APPOINTMENTS): self.appointments = dict() self.locator_uuid_map = dict() self.asleep = True self.block_queue = Queue() self.max_appointments = max_appointments - self.zmq_subscriber = None + self.chain_monitor = chain_monitor self.db_manager = db_manager self.signing_key = Cryptographer.load_private_key_der(sk_der) @@ -127,10 +126,8 @@ class Watcher: if self.asleep: self.asleep = False - zmq_thread = Thread(target=self.do_subscribe) - watcher = Thread(target=self.do_watch) - zmq_thread.start() - watcher.start() + self.chain_monitor.watcher_asleep = False + Thread(target=self.do_watch).start() logger.info("Waking up") @@ -151,15 +148,6 @@ class Watcher: return appointment_added, signature - def do_subscribe(self): - """ - Initializes a ``ZMQSubscriber`` instance to listen to new blocks from ``bitcoind``. Block ids are received - trough the ``block_queue``. - """ - - self.zmq_subscriber = ZMQSubscriber(parent="Watcher") - self.zmq_subscriber.handle(self.block_queue) - def do_watch(self): """ Monitors the blockchain whilst there are pending appointments. @@ -221,8 +209,7 @@ class Watcher: # Go back to sleep if there are no more appointments self.asleep = True - self.zmq_subscriber.terminate = True - self.block_queue = Queue() + self.chain_monitor.watcher_asleep = True logger.info("No more pending appointments, going back to sleep") From c594018dce377968cb212b10dcdf2c157fd9f700 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Fri, 3 Jan 2020 11:37:26 +0100 Subject: [PATCH 02/14] Bug fix: Properly initializes Responder --- pisa/chain_monitor.py | 3 +-- pisa/pisad.py | 13 ++++++------- pisa/watcher.py | 2 +- 3 files changed, 8 insertions(+), 10 deletions(-) diff --git a/pisa/chain_monitor.py b/pisa/chain_monitor.py index 74b21f1..4442488 100644 --- a/pisa/chain_monitor.py +++ b/pisa/chain_monitor.py @@ -1,6 +1,5 @@ import zmq import binascii -from queue import Queue from threading import Thread, Event, Condition from common.logger import Logger @@ -26,8 +25,8 @@ class ChainMonitor: self.zmqSubSocket.connect("%s://%s:%s" % (FEED_PROTOCOL, FEED_ADDR, FEED_PORT)) self.watcher_queue = None - self.watcher_asleep = True self.responder_queue = None + self.watcher_asleep = True self.responder_asleep = True def attach_watcher(self, queue, awake): diff --git a/pisa/pisad.py b/pisa/pisad.py index ed185c6..4935d40 100644 --- a/pisa/pisad.py +++ b/pisa/pisad.py @@ -71,7 +71,6 @@ if __name__ == "__main__": last_block_responder = db_manager.load_last_block_hash_responder() # FIXME: 32-reorgs-offline dropped txs are not used at this point. - responder = Responder(db_manager, chain_monitor) last_common_ancestor_responder = None missed_blocks_responder = None @@ -82,13 +81,13 @@ if __name__ == "__main__": ) missed_blocks_responder = block_processor.get_missed_blocks(last_common_ancestor_responder) - responder.trackers, responder.tx_tracker_map = Builder.build_trackers(responder_trackers_data) - responder.block_queue = Builder.build_block_queue(missed_blocks_responder) + watcher.responder.trackers, watcher.responder.tx_tracker_map = Builder.build_trackers( + responder_trackers_data + ) + watcher.responder.block_queue = Builder.build_block_queue(missed_blocks_responder) - # Build Watcher with Responder and backed up data. If the blocks of both match we don't perform the - # search twice. - watcher.responder = responder - chain_monitor.attach_responder(responder.block_queue, responder.asleep) + # Build Watcher. If the blocks of both match we don't perform the search twice. + chain_monitor.attach_responder(watcher.responder.block_queue, watcher.responder.asleep) if last_block_watcher is not None: if last_block_watcher == last_block_responder: diff --git a/pisa/watcher.py b/pisa/watcher.py index d6e0729..1da5c75 100644 --- a/pisa/watcher.py +++ b/pisa/watcher.py @@ -68,7 +68,7 @@ class Watcher: self.signing_key = Cryptographer.load_private_key_der(sk_der) if not isinstance(responder, Responder): - self.responder = Responder(db_manager) + self.responder = Responder(db_manager, chain_monitor) @staticmethod def compute_locator(tx_id): From 3889d30e313c75dace4e321b25282e4fcf68830b Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Fri, 3 Jan 2020 13:35:52 +0100 Subject: [PATCH 03/14] Minor refactor Runs chain_monitor as daemon so it terminates along with the main thread. Reorders attack_responder so it is in the same stanza as attack_watcher --- pisa/chain_monitor.py | 6 +++--- pisa/pisad.py | 3 +-- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/pisa/chain_monitor.py b/pisa/chain_monitor.py index 4442488..c9e9f83 100644 --- a/pisa/chain_monitor.py +++ b/pisa/chain_monitor.py @@ -53,11 +53,11 @@ class ChainMonitor: def monitor_chain(self): self.best_tip = BlockProcessor.get_best_block_hash() - Thread(target=self.monitor_chain_polling).start() - Thread(target=self.monitor_chain_zmq).start() + Thread(target=self.monitor_chain_polling, daemon=True).start() + Thread(target=self.monitor_chain_zmq, daemon=True).start() def monitor_chain_polling(self): - while self.terminate: + while not self.terminate: self.check_tip.wait(timeout=60) # Terminate could have been set wile the thread was blocked in wait diff --git a/pisa/pisad.py b/pisa/pisad.py index 4935d40..b027662 100644 --- a/pisa/pisad.py +++ b/pisa/pisad.py @@ -59,6 +59,7 @@ if __name__ == "__main__": watcher = Watcher(db_manager, chain_monitor, secret_key_der) chain_monitor.attach_watcher(watcher.block_queue, watcher.asleep) + chain_monitor.attach_responder(watcher.responder.block_queue, watcher.responder.asleep) if len(watcher_appointments_data) == 0 and len(responder_trackers_data) == 0: logger.info("Fresh bootstrap") @@ -87,8 +88,6 @@ if __name__ == "__main__": watcher.responder.block_queue = Builder.build_block_queue(missed_blocks_responder) # Build Watcher. If the blocks of both match we don't perform the search twice. - chain_monitor.attach_responder(watcher.responder.block_queue, watcher.responder.asleep) - if last_block_watcher is not None: if last_block_watcher == last_block_responder: missed_blocks_watcher = missed_blocks_responder From 069e15fdba0ab0166adaf23b3e4327fa50064feb Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Fri, 3 Jan 2020 13:36:52 +0100 Subject: [PATCH 04/14] Updates current tests to work with chain_monitor instead of zmq_sub --- test/pisa/unit/test_api.py | 9 ++- test/pisa/unit/test_responder.py | 106 ++++++++++++++----------------- test/pisa/unit/test_watcher.py | 62 ++++++++---------- 3 files changed, 81 insertions(+), 96 deletions(-) diff --git a/test/pisa/unit/test_api.py b/test/pisa/unit/test_api.py index 2dc830b..b6ca968 100644 --- a/test/pisa/unit/test_api.py +++ b/test/pisa/unit/test_api.py @@ -10,6 +10,7 @@ from pisa.watcher import Watcher from pisa.tools import bitcoin_cli from pisa import HOST, PORT from pisa.conf import MAX_APPOINTMENTS +from pisa.chain_monitor import ChainMonitor from test.pisa.unit.conftest import ( generate_block, @@ -37,7 +38,13 @@ def run_api(db_manager): format=serialization.PrivateFormat.TraditionalOpenSSL, encryption_algorithm=serialization.NoEncryption(), ) - watcher = Watcher(db_manager, sk_der) + + chain_monitor = ChainMonitor() + chain_monitor.monitor_chain() + + watcher = Watcher(db_manager, chain_monitor, sk_der) + chain_monitor.attach_watcher(watcher.block_queue, watcher.asleep) + chain_monitor.attach_responder(watcher.responder.block_queue, watcher.responder.asleep) api_thread = Thread(target=API(watcher).start) api_thread.daemon = True diff --git a/test/pisa/unit/test_responder.py b/test/pisa/unit/test_responder.py index a9e99da..b3e339e 100644 --- a/test/pisa/unit/test_responder.py +++ b/test/pisa/unit/test_responder.py @@ -5,15 +5,14 @@ from uuid import uuid4 from shutil import rmtree from copy import deepcopy from threading import Thread -from queue import Queue, Empty from pisa.db_manager import DBManager from pisa.responder import Responder, TransactionTracker from pisa.block_processor import BlockProcessor +from pisa.chain_monitor import ChainMonitor from pisa.tools import bitcoin_cli from common.constants import LOCATOR_LEN_HEX -from common.tools import check_sha256_hex_format from bitcoind_mock.utils import sha256d from bitcoind_mock.transaction import TX @@ -21,8 +20,19 @@ from test.pisa.unit.conftest import generate_block, generate_blocks, get_random_ @pytest.fixture(scope="module") -def responder(db_manager): - return Responder(db_manager) +def chain_monitor(): + chain_monitor = ChainMonitor() + chain_monitor.monitor_chain() + + return chain_monitor + + +@pytest.fixture(scope="module") +def responder(db_manager, chain_monitor): + responder = Responder(db_manager, chain_monitor) + chain_monitor.attach_responder(responder.block_queue, responder.asleep) + + return responder @pytest.fixture() @@ -145,17 +155,19 @@ def test_tracker_from_dict_invalid_data(): def test_init_responder(responder): - assert type(responder.trackers) is dict and len(responder.trackers) == 0 - assert type(responder.tx_tracker_map) is dict and len(responder.tx_tracker_map) == 0 - assert type(responder.unconfirmed_txs) is list and len(responder.unconfirmed_txs) == 0 - assert type(responder.missed_confirmations) is dict and len(responder.missed_confirmations) == 0 + assert isinstance(responder.trackers, dict) and len(responder.trackers) == 0 + assert isinstance(responder.tx_tracker_map, dict) and len(responder.tx_tracker_map) == 0 + assert isinstance(responder.unconfirmed_txs, list) and len(responder.unconfirmed_txs) == 0 + assert isinstance(responder.missed_confirmations, dict) and len(responder.missed_confirmations) == 0 + assert isinstance(responder.chain_monitor, ChainMonitor) assert responder.block_queue.empty() assert responder.asleep is True - assert responder.zmq_subscriber is None -def test_handle_breach(db_manager): - responder = Responder(db_manager) +def test_handle_breach(db_manager, chain_monitor): + responder = Responder(db_manager, chain_monitor) + chain_monitor.attach_responder(responder.block_queue, responder.asleep) + uuid = uuid4().hex tracker = create_dummy_tracker() @@ -172,11 +184,10 @@ def test_handle_breach(db_manager): assert receipt.delivered is True - # The responder automatically fires add_tracker on adding a tracker if it is asleep. We need to stop the processes now. - # To do so we delete all the trackers, stop the zmq and create a new fake block to unblock the queue.get method + # The responder automatically fires add_tracker on adding a tracker if it is asleep. We need to stop the processes + # now. To do so we delete all the trackers, and generate a new block. responder.trackers = dict() - responder.zmq_subscriber.terminate = True - responder.block_queue.put(get_random_value_hex(32)) + generate_block() def test_add_bad_response(responder): @@ -205,7 +216,7 @@ def test_add_bad_response(responder): def test_add_tracker(responder): - responder.asleep = False + # Responder is asleep for _ in range(20): uuid = uuid4().hex @@ -237,7 +248,8 @@ def test_add_tracker(responder): def test_add_tracker_same_penalty_txid(responder): - # Create the same tracker using two different uuids + # Responder is asleep + confirmations = 0 locator, dispute_txid, penalty_txid, penalty_rawtx, appointment_end = create_dummy_tracker_data(random_txid=True) uuid_1 = uuid4().hex @@ -264,7 +276,7 @@ def test_add_tracker_same_penalty_txid(responder): def test_add_tracker_already_confirmed(responder): - responder.asleep = False + # Responder is asleep for i in range(20): uuid = uuid4().hex @@ -278,29 +290,10 @@ def test_add_tracker_already_confirmed(responder): assert penalty_txid not in responder.unconfirmed_txs -def test_do_subscribe(responder): - responder.block_queue = Queue() - - zmq_thread = Thread(target=responder.do_subscribe) - zmq_thread.daemon = True - zmq_thread.start() - - try: - generate_block() - block_hash = responder.block_queue.get() - assert check_sha256_hex_format(block_hash) - - except Empty: - assert False - - -def test_do_watch(temp_db_manager): - responder = Responder(temp_db_manager) - responder.block_queue = Queue() - - zmq_thread = Thread(target=responder.do_subscribe) - zmq_thread.daemon = True - zmq_thread.start() +def test_do_watch(temp_db_manager, chain_monitor): + # Create a fresh responder to simplify the test + responder = Responder(temp_db_manager, chain_monitor) + chain_monitor.attach_responder(responder.block_queue, False) trackers = [create_dummy_tracker(penalty_rawtx=TX.create_dummy_transaction()) for _ in range(20)] @@ -314,9 +307,7 @@ def test_do_watch(temp_db_manager): responder.unconfirmed_txs.append(tracker.penalty_txid) # Let's start to watch - watch_thread = Thread(target=responder.do_watch) - watch_thread.daemon = True - watch_thread.start() + Thread(target=responder.do_watch, daemon=True).start() # And broadcast some of the transactions broadcast_txs = [] @@ -324,8 +315,10 @@ def test_do_watch(temp_db_manager): bitcoin_cli().sendrawtransaction(tracker.penalty_rawtx) broadcast_txs.append(tracker.penalty_txid) + print(responder.unconfirmed_txs) # Mine a block generate_block() + print(responder.unconfirmed_txs) # The transactions we sent shouldn't be in the unconfirmed transaction list anymore assert not set(broadcast_txs).issubset(responder.unconfirmed_txs) @@ -350,13 +343,9 @@ def test_do_watch(temp_db_manager): assert responder.asleep is True -def test_check_confirmations(temp_db_manager): - responder = Responder(temp_db_manager) - responder.block_queue = Queue() - - zmq_thread = Thread(target=responder.do_subscribe) - zmq_thread.daemon = True - zmq_thread.start() +def test_check_confirmations(temp_db_manager, chain_monitor): + responder = Responder(temp_db_manager, chain_monitor) + chain_monitor.attach_responder(responder.block_queue, responder.asleep) # check_confirmations checks, given a list of transaction for a block, what of the known penalty transaction have # been confirmed. To test this we need to create a list of transactions and the state of the responder @@ -386,7 +375,7 @@ def test_check_confirmations(temp_db_manager): assert responder.missed_confirmations[tx] == 1 -# WIP: Check this properly, a bug pass unnoticed! +# TODO: Check this properly, a bug pass unnoticed! def test_get_txs_to_rebroadcast(responder): # Let's create a few fake txids and assign at least 6 missing confirmations to each txs_missing_too_many_conf = {get_random_value_hex(32): 6 + i for i in range(10)} @@ -410,13 +399,13 @@ def test_get_txs_to_rebroadcast(responder): assert txs_to_rebroadcast == list(txs_missing_too_many_conf.keys()) -def test_get_completed_trackers(db_manager): +def test_get_completed_trackers(db_manager, chain_monitor): initial_height = bitcoin_cli().getblockcount() - # Let's use a fresh responder for this to make it easier to compare the results - responder = Responder(db_manager) + responder = Responder(db_manager, chain_monitor) + chain_monitor.attach_responder(responder.block_queue, responder.asleep) - # A complete tracker is a tracker that has reached the appointment end with enough confirmations (> MIN_CONFIRMATIONS) + # A complete tracker is a tracker that has reached the appointment end with enough confs (> MIN_CONFIRMATIONS) # We'll create three type of transactions: end reached + enough conf, end reached + no enough conf, end not reached trackers_end_conf = { uuid4().hex: create_dummy_tracker(penalty_rawtx=TX.create_dummy_transaction()) for _ in range(10) @@ -461,9 +450,10 @@ def test_get_completed_trackers(db_manager): assert set(completed_trackers_ids) == set(ended_trackers_keys) -def test_rebroadcast(db_manager): - responder = Responder(db_manager) +def test_rebroadcast(db_manager, chain_monitor): + responder = Responder(db_manager, chain_monitor) responder.asleep = False + chain_monitor.attach_responder(responder.block_queue, responder.asleep) txs_to_rebroadcast = [] diff --git a/test/pisa/unit/test_watcher.py b/test/pisa/unit/test_watcher.py index 8ec331c..2b2c0cb 100644 --- a/test/pisa/unit/test_watcher.py +++ b/test/pisa/unit/test_watcher.py @@ -1,22 +1,15 @@ import pytest from uuid import uuid4 from threading import Thread -from queue import Queue, Empty from cryptography.hazmat.primitives import serialization from pisa.watcher import Watcher from pisa.responder import Responder from pisa.tools import bitcoin_cli -from test.pisa.unit.conftest import ( - generate_block, - generate_blocks, - generate_dummy_appointment, - get_random_value_hex, - generate_keypair, -) +from test.pisa.unit.conftest import generate_blocks, generate_dummy_appointment, get_random_value_hex, generate_keypair +from pisa.chain_monitor import ChainMonitor from pisa.conf import EXPIRY_DELTA, MAX_APPOINTMENTS -from common.tools import check_sha256_hex_format from common.cryptographer import Cryptographer @@ -35,8 +28,20 @@ sk_der = signing_key.private_bytes( @pytest.fixture(scope="module") -def watcher(db_manager): - return Watcher(db_manager, sk_der) +def chain_monitor(): + chain_monitor = ChainMonitor() + chain_monitor.monitor_chain() + + return chain_monitor + + +@pytest.fixture(scope="module") +def watcher(db_manager, chain_monitor): + watcher = Watcher(db_manager, chain_monitor, sk_der) + chain_monitor.attach_watcher(watcher.block_queue, watcher.asleep) + chain_monitor.attach_responder(watcher.responder.block_queue, watcher.responder.asleep) + + return watcher @pytest.fixture(scope="module") @@ -67,17 +72,18 @@ def create_appointments(n): return appointments, locator_uuid_map, dispute_txs -def test_init(watcher): - assert type(watcher.appointments) is dict and len(watcher.appointments) == 0 - assert type(watcher.locator_uuid_map) is dict and len(watcher.locator_uuid_map) == 0 +def test_init(run_bitcoind, watcher): + assert isinstance(watcher.appointments, dict) and len(watcher.appointments) == 0 + assert isinstance(watcher.locator_uuid_map, dict) and len(watcher.locator_uuid_map) == 0 + assert isinstance(watcher.chain_monitor, ChainMonitor) assert watcher.block_queue.empty() assert watcher.asleep is True + assert watcher.max_appointments == MAX_APPOINTMENTS - assert watcher.zmq_subscriber is None assert type(watcher.responder) is Responder -def test_add_appointment(run_bitcoind, watcher): +def test_add_appointment(watcher): # The watcher automatically fires do_watch and do_subscribe on adding an appointment if it is asleep (initial state) # Avoid this by setting the state to awake. watcher.asleep = False @@ -121,36 +127,18 @@ def test_add_too_many_appointments(watcher): assert sig is None -def test_do_subscribe(watcher): - watcher.block_queue = Queue() - - zmq_thread = Thread(target=watcher.do_subscribe) - zmq_thread.daemon = True - zmq_thread.start() - - try: - generate_block() - block_hash = watcher.block_queue.get() - assert check_sha256_hex_format(block_hash) - - except Empty: - assert False - - def test_do_watch(watcher): # We will wipe all the previous data and add 5 appointments watcher.appointments, watcher.locator_uuid_map, dispute_txs = create_appointments(APPOINTMENTS) + watcher.chain_monitor.watcher_asleep = False - watch_thread = Thread(target=watcher.do_watch) - watch_thread.daemon = True - watch_thread.start() + Thread(target=watcher.do_watch, daemon=True).start() # Broadcast the first two for dispute_tx in dispute_txs[:2]: bitcoin_cli().sendrawtransaction(dispute_tx) - # After leaving some time for the block to be mined and processed, the number of appointments should have reduced - # by two + # After generating enough blocks, the number of appointments should have reduced by two generate_blocks(START_TIME_OFFSET + END_TIME_OFFSET) assert len(watcher.appointments) == APPOINTMENTS - 2 From e5514cefcec439d3a84a45402562c0754f86d275 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Mon, 6 Jan 2020 12:27:02 +0100 Subject: [PATCH 05/14] Updates old docs/comments regarding zmq --- pisa/responder.py | 9 ++++----- pisa/watcher.py | 8 ++++---- test/pisa/unit/test_responder.py | 2 +- 3 files changed, 9 insertions(+), 10 deletions(-) diff --git a/pisa/responder.py b/pisa/responder.py index dbfac15..89b1096 100644 --- a/pisa/responder.py +++ b/pisa/responder.py @@ -126,9 +126,9 @@ class Responder: has missed. Used to trigger rebroadcast if needed. asleep (:obj:`bool`): A flag that signals whether the :obj:`Responder` is asleep or awake. block_queue (:obj:`Queue`): A queue used by the :obj:`Responder` to receive block hashes from ``bitcoind``. It - is populated by the :obj:`ZMQSubscriber `. - zmq_subscriber (:obj:`ZMQSubscriber `): a ``ZMQSubscriber`` instance - used to receive new block notifications from ``bitcoind``. + is populated by the :obj:`ChainMonitor `. + chain_monitor (:obj:`ChainMonitor `): a ``ChainMonitor`` instance used to track + new blocks received by ``bitcoind``. db_manager (:obj:`DBManager `): A ``DBManager`` instance to interact with the database. @@ -223,8 +223,7 @@ class Responder: The :obj:`TransactionTracker` is stored in ``trackers`` and ``tx_tracker_map`` and the ``penalty_txid`` added to ``unconfirmed_txs`` if ``confirmations=0``. Finally, the data is also stored in the database. - ``add_tracker`` awakes the :obj:`Responder` and creates a connection with the - :obj:`ZMQSubscriber ` if he is asleep. + ``add_tracker`` awakes the :obj:`Responder` if it is asleep. Args: uuid (:obj:`str`): a unique identifier for the appointment. diff --git a/pisa/watcher.py b/pisa/watcher.py index 1da5c75..079953b 100644 --- a/pisa/watcher.py +++ b/pisa/watcher.py @@ -27,7 +27,7 @@ class Watcher: If an appointment reaches its end with no breach, the data is simply deleted. The :class:`Watcher` receives information about new received blocks via the ``block_queue`` that is populated by the - :obj:`ZMQSubscriber `. + :obj:`ChainMonitor `. Args: db_manager (:obj:`DBManager `): a ``DBManager`` instance to interact with the database. @@ -45,11 +45,11 @@ class Watcher: appointments with the same ``locator``. asleep (:obj:`bool`): A flag that signals whether the :obj:`Watcher` is asleep or awake. block_queue (:obj:`Queue`): A queue used by the :obj:`Watcher` to receive block hashes from ``bitcoind``. It is - populated by the :obj:`ZMQSubscriber `. + populated by the :obj:`ChainMonitor `. max_appointments(:obj:`int`): the maximum amount of appointments that the :obj:`Watcher` will keep at any given time. - zmq_subscriber (:obj:`ZMQSubscriber `): a ZMQSubscriber instance used - to receive new block notifications from ``bitcoind``. + chain_monitor (:obj:`ChainMonitor `): a ``ChainMonitor`` instance used to track + new blocks received by ``bitcoind``. db_manager (:obj:`DBManager `): A db manager instance to interact with the database. Raises: diff --git a/test/pisa/unit/test_responder.py b/test/pisa/unit/test_responder.py index b3e339e..d515ff7 100644 --- a/test/pisa/unit/test_responder.py +++ b/test/pisa/unit/test_responder.py @@ -195,7 +195,7 @@ def test_add_bad_response(responder): tracker = create_dummy_tracker() # Now that the asleep / awake functionality has been tested we can avoid manually killing the responder by setting - # to awake. That will prevent the zmq thread to be launched again. + # to awake. That will prevent the chain_monitor thread to be launched again. responder.asleep = False # A txid instead of a rawtx should be enough for unit tests using the bitcoind mock, better tests are needed though. From b2272aa4eaeec2c19f0cc2385b9593da60826636 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Mon, 6 Jan 2020 13:32:25 +0100 Subject: [PATCH 06/14] Adds ChainMonitor config parameters --- pisa/sample_conf.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/pisa/sample_conf.py b/pisa/sample_conf.py index 71e64a0..8d08590 100644 --- a/pisa/sample_conf.py +++ b/pisa/sample_conf.py @@ -5,6 +5,10 @@ BTC_RPC_HOST = "localhost" BTC_RPC_PORT = 18443 BTC_NETWORK = "regtest" +# CHAIN MONITOR +POLLING_DELTA = 60 +BLOCK_WINDOW_SIZE = 10 + # ZMQ FEED_PROTOCOL = "tcp" FEED_ADDR = "127.0.0.1" From 9bb69d1f5a51753b5828040f4aa2e25d282dc525 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Mon, 6 Jan 2020 13:32:35 +0100 Subject: [PATCH 07/14] Uses conf parameters instead of hardcoded values for the ChainMonitor and includes docstrings --- pisa/chain_monitor.py | 106 ++++++++++++++++++++++++++++++++++++------ 1 file changed, 92 insertions(+), 14 deletions(-) diff --git a/pisa/chain_monitor.py b/pisa/chain_monitor.py index c9e9f83..fdbc2b4 100644 --- a/pisa/chain_monitor.py +++ b/pisa/chain_monitor.py @@ -3,13 +3,36 @@ import binascii from threading import Thread, Event, Condition from common.logger import Logger -from pisa.conf import FEED_PROTOCOL, FEED_ADDR, FEED_PORT +from pisa.conf import FEED_PROTOCOL, FEED_ADDR, FEED_PORT, POLLING_DELTA, BLOCK_WINDOW_SIZE from pisa.block_processor import BlockProcessor logger = Logger("ChainMonitor") class ChainMonitor: + """ + The :class:`ChainMonitor` is the class in charge of monitoring the blockchain (via ``bitcoind``) to detect new + blocks on top of the best chain. If a new best block is spotted, the chain monitor will notify the + :obj:`Watcher ` and the :obj:`Responder ` using ``Queues``. + + The :class:`ChainMonitor` monitors the chain using two methods: ``zmq`` and ``polling``. Blocks are only notified + once per queue and the notification is triggered by the method that detects the block faster. + + Attributes: + best_tip (:obj:`str`): a block hash representing the current best tip. + last_tips (:obj:`list`): a list of last chain tips. Used as an sliding window to avoid notifying about old tips. + terminate (:obj:`bool`): a flag to signal the termination of the :class:`ChainMonitor` (shutdown the tower). + check_tip (:obj:`Event`): an event that it's triggered at fixed time intervals and controls the polling thread. + lock (:obj:`Condition`): a lock used to protect concurrent access to the queues and ``best_tip`` by the zmq and + polling threads. + zmqSubSocket (:obj:`socket`): a socket to connect to ``bitcoind`` via ``zmq``. + watcher_queue (:obj:`Queue`): a queue to send new best tips to the :obj:`Watcher `. + responder_queue (:obj:`Queue`): a queue to send new best tips to the + :obj:`Responder `. + watcher_asleep (:obj:`bool`): a flag that signals whether to send information to the ``Watcher`` or not. + responder_asleep (:obj:`bool`): a flag that signals whether to send information to the ``Responder`` or not. + """ + def __init__(self): self.best_tip = None self.last_tips = [] @@ -29,36 +52,74 @@ class ChainMonitor: self.watcher_asleep = True self.responder_asleep = True - def attach_watcher(self, queue, awake): - self.watcher_queue = queue - self.watcher_asleep = awake + def attach_watcher(self, queue, asleep): + """ + Attaches a :obj:`Watcher ` to the :class:`ChainMonitor`. The ``Watcher`` and the + ``ChainMonitor`` are connected via the ``watcher_queue`` and the ``watcher_asleep`` flag. + + Args: + queue (:obj:`Queue`): the queue to be used to send blocks hashes to the ``Watcher``. + asleep( :obj:`bool`): whether the ``Watcher`` is initially awake or asleep. It is changed on the fly from + the ``Watcher`` when the state changes. + """ + + self.watcher_queue = queue + self.watcher_asleep = asleep + + def attach_responder(self, queue, asleep): + """ + Attaches a :obj:`Responder ` to the :class:`ChainMonitor`. The ``Responder`` and the + ``ChainMonitor`` are connected via the ``responder_queue`` and the ``responder_asleep`` flag. + + Args: + queue (:obj:`Queue`): the queue to be used to send blocks hashes to the ``Responder``. + asleep( :obj:`bool`): whether the ``Responder`` is initially awake or asleep. It is changed on the fly from + the ``Responder`` when the state changes. + """ - def attach_responder(self, queue, awake): self.responder_queue = queue - self.responder_asleep = awake + self.responder_asleep = asleep def notify_subscribers(self, block_hash): + """ + Notifies the subscribers (``Watcher`` and ``Responder``) about a new block provided they are awake. It does so + by putting the hash in the corresponding queue(s). + + Args: + block_hash (:obj:`str`): the new block hash to be sent to the awake subscribers. + """ + if not self.watcher_asleep: self.watcher_queue.put(block_hash) if not self.responder_asleep: self.responder_queue.put(block_hash) - def update_state(self, block_hash): + def update_state(self, block_hash, max_block_window_size=BLOCK_WINDOW_SIZE): + """ + Updates the state of the ``ChainMonitor``. The state is represented as the ``best_tip`` and the list of + ``last_tips``. ``last_tips`` is bounded to ``max_block_window_size``. + + Args: + block_hash (:obj:`block_hash`): the new best tip. + max_block_window_size (:obj:`int`): the maximum length of the ``last_tips`` list. + """ + self.best_tip = block_hash self.last_tips.append(block_hash) - if len(self.last_tips) > 10: + if len(self.last_tips) > max_block_window_size: self.last_tips.pop(0) - def monitor_chain(self): - self.best_tip = BlockProcessor.get_best_block_hash() - Thread(target=self.monitor_chain_polling, daemon=True).start() - Thread(target=self.monitor_chain_zmq, daemon=True).start() + def monitor_chain_polling(self, polling_delta=POLLING_DELTA): + """ + Monitors ``bitcoind`` via polling. Once the method is fired, it keeps monitoring as long as ``terminate`` is not + set. Polling is performed once every ``polling_delta`` seconds. If a new best tip if found, the shared lock is + acquired, the state is updated and the subscribers are notified, and finally the lock is released. + """ - def monitor_chain_polling(self): while not self.terminate: - self.check_tip.wait(timeout=60) + self.check_tip.wait(timeout=polling_delta) # Terminate could have been set wile the thread was blocked in wait if not self.terminate: @@ -72,6 +133,12 @@ class ChainMonitor: self.lock.release() def monitor_chain_zmq(self): + """ + Monitors ``bitcoind`` via zmq. Once the method is fired, it keeps monitoring as long as ``terminate`` is not + set. If a new best tip if found, the shared lock is acquired, the state is updated and the subscribers are + notified, and finally the lock is released. + """ + while not self.terminate: msg = self.zmqSubSocket.recv_multipart() @@ -89,3 +156,14 @@ class ChainMonitor: self.notify_subscribers(block_hash) logger.info("New block received via zmq", block_hash=block_hash) self.lock.release() + + def monitor_chain(self): + """ + Main :class:`ChainMonitor` method. It initializes the ``best_tip`` to the current one (by querying the + :obj:`BlockProcessor `) and creates two threads, one per each monitoring + approach (``zmq`` and ``polling``) + """ + + self.best_tip = BlockProcessor.get_best_block_hash() + Thread(target=self.monitor_chain_polling, daemon=True).start() + Thread(target=self.monitor_chain_zmq, daemon=True).start() From f10c3c46eb6c96e8cc625db2f8d3e4ac36470c2d Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Tue, 7 Jan 2020 16:08:10 +0100 Subject: [PATCH 08/14] Updates tests to use ChainMonitor as global fixture --- test/pisa/unit/conftest.py | 12 ++++++++++++ test/pisa/unit/test_api.py | 5 +---- test/pisa/unit/test_responder.py | 8 -------- test/pisa/unit/test_watcher.py | 8 -------- 4 files changed, 13 insertions(+), 20 deletions(-) diff --git a/test/pisa/unit/conftest.py b/test/pisa/unit/conftest.py index 4ff9028..6616611 100644 --- a/test/pisa/unit/conftest.py +++ b/test/pisa/unit/conftest.py @@ -15,6 +15,7 @@ from pisa.responder import TransactionTracker from pisa.watcher import Watcher from pisa.tools import bitcoin_cli from pisa.db_manager import DBManager +from pisa.chain_monitor import ChainMonitor from common.appointment import Appointment from bitcoind_mock.utils import sha256d @@ -50,6 +51,17 @@ def db_manager(): rmtree("test_db") +@pytest.fixture(scope="module") +def chain_monitor(): + chain_monitor = ChainMonitor() + chain_monitor.monitor_chain() + + yield chain_monitor + + chain_monitor.terminate = True + generate_block() + + def generate_keypair(): client_sk = ec.generate_private_key(ec.SECP256K1, default_backend()) client_pk = client_sk.public_key() diff --git a/test/pisa/unit/test_api.py b/test/pisa/unit/test_api.py index b6ca968..1e1d4ef 100644 --- a/test/pisa/unit/test_api.py +++ b/test/pisa/unit/test_api.py @@ -31,7 +31,7 @@ locator_dispute_tx_map = {} @pytest.fixture(scope="module") -def run_api(db_manager): +def run_api(db_manager, chain_monitor): sk, pk = generate_keypair() sk_der = sk.private_bytes( encoding=serialization.Encoding.DER, @@ -39,9 +39,6 @@ def run_api(db_manager): encryption_algorithm=serialization.NoEncryption(), ) - chain_monitor = ChainMonitor() - chain_monitor.monitor_chain() - watcher = Watcher(db_manager, chain_monitor, sk_der) chain_monitor.attach_watcher(watcher.block_queue, watcher.asleep) chain_monitor.attach_responder(watcher.responder.block_queue, watcher.responder.asleep) diff --git a/test/pisa/unit/test_responder.py b/test/pisa/unit/test_responder.py index d515ff7..b5f5949 100644 --- a/test/pisa/unit/test_responder.py +++ b/test/pisa/unit/test_responder.py @@ -19,14 +19,6 @@ from bitcoind_mock.transaction import TX from test.pisa.unit.conftest import generate_block, generate_blocks, get_random_value_hex -@pytest.fixture(scope="module") -def chain_monitor(): - chain_monitor = ChainMonitor() - chain_monitor.monitor_chain() - - return chain_monitor - - @pytest.fixture(scope="module") def responder(db_manager, chain_monitor): responder = Responder(db_manager, chain_monitor) diff --git a/test/pisa/unit/test_watcher.py b/test/pisa/unit/test_watcher.py index 2b2c0cb..80dd396 100644 --- a/test/pisa/unit/test_watcher.py +++ b/test/pisa/unit/test_watcher.py @@ -27,14 +27,6 @@ sk_der = signing_key.private_bytes( ) -@pytest.fixture(scope="module") -def chain_monitor(): - chain_monitor = ChainMonitor() - chain_monitor.monitor_chain() - - return chain_monitor - - @pytest.fixture(scope="module") def watcher(db_manager, chain_monitor): watcher = Watcher(db_manager, chain_monitor, sk_der) From 31b4c2e993c8e33a7e281b3d748ea953cf06ca07 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Tue, 7 Jan 2020 16:08:43 +0100 Subject: [PATCH 09/14] Bug fix and moves update condition to update_state The tip update was performed in the wrong order, so tips never made it to last_tips --- pisa/chain_monitor.py | 23 +++++++++++++++-------- 1 file changed, 15 insertions(+), 8 deletions(-) diff --git a/pisa/chain_monitor.py b/pisa/chain_monitor.py index fdbc2b4..087096b 100644 --- a/pisa/chain_monitor.py +++ b/pisa/chain_monitor.py @@ -103,13 +103,22 @@ class ChainMonitor: Args: block_hash (:obj:`block_hash`): the new best tip. max_block_window_size (:obj:`int`): the maximum length of the ``last_tips`` list. + + Returns: + (:obj:`bool`): ``True`` is the state was successfully updated, ``False`` otherwise. """ - self.best_tip = block_hash - self.last_tips.append(block_hash) + if block_hash != self.best_tip and block_hash not in self.last_tips: + self.last_tips.append(self.best_tip) + self.best_tip = block_hash - if len(self.last_tips) > max_block_window_size: - self.last_tips.pop(0) + if len(self.last_tips) > max_block_window_size: + self.last_tips.pop(0) + + return True + + else: + return False def monitor_chain_polling(self, polling_delta=POLLING_DELTA): """ @@ -126,8 +135,7 @@ class ChainMonitor: current_tip = BlockProcessor.get_best_block_hash() self.lock.acquire() - if current_tip != self.best_tip: - self.update_state(current_tip) + if self.update_state(current_tip): self.notify_subscribers(current_tip) logger.info("New block received via polling", block_hash=current_tip) self.lock.release() @@ -151,8 +159,7 @@ class ChainMonitor: block_hash = binascii.hexlify(body).decode("utf-8") self.lock.acquire() - if block_hash != self.best_tip and block_hash not in self.last_tips: - self.update_state(block_hash) + if self.update_state(block_hash): self.notify_subscribers(block_hash) logger.info("New block received via zmq", block_hash=block_hash) self.lock.release() From 40e656dcd3571f8aea52a33a2c9aa69f1afe5041 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Tue, 7 Jan 2020 16:09:49 +0100 Subject: [PATCH 10/14] ChainMonitor unit tests --- test/pisa/unit/test_chain_monitor.py | 194 +++++++++++++++++++++++++++ 1 file changed, 194 insertions(+) create mode 100644 test/pisa/unit/test_chain_monitor.py diff --git a/test/pisa/unit/test_chain_monitor.py b/test/pisa/unit/test_chain_monitor.py new file mode 100644 index 0000000..d3037d8 --- /dev/null +++ b/test/pisa/unit/test_chain_monitor.py @@ -0,0 +1,194 @@ +import zmq +import time +from threading import Thread, Event, Condition + +from pisa.watcher import Watcher +from pisa.responder import Responder +from pisa.block_processor import BlockProcessor +from pisa.chain_monitor import ChainMonitor + +from test.pisa.unit.conftest import get_random_value_hex, generate_block + + +def test_init(run_bitcoind): + # run_bitcoind is started here instead of later on to avoid race conditions while it initializes + + # Not much to test here, just sanity checks to make sure nothings goes south in the future + chain_monitor = ChainMonitor() + + assert chain_monitor.best_tip is None + assert isinstance(chain_monitor.last_tips, list) and len(chain_monitor.last_tips) == 0 + assert chain_monitor.terminate is False + assert isinstance(chain_monitor.check_tip, Event) + assert isinstance(chain_monitor.lock, Condition) + assert isinstance(chain_monitor.zmqSubSocket, zmq.Socket) + + # The Queues and asleep flags are initialized when attaching the corresponding subscriber + assert chain_monitor.watcher_queue is None + assert chain_monitor.responder_queue is None + assert chain_monitor.watcher_asleep and chain_monitor.responder_asleep + + +def test_attach_watcher(chain_monitor): + watcher = Watcher(db_manager=None, chain_monitor=chain_monitor, sk_der=None) + chain_monitor.attach_watcher(watcher.block_queue, watcher.asleep) + + # booleans are not passed as reference in Python, so the flags need to be set separately + assert watcher.asleep == chain_monitor.watcher_asleep + watcher.asleep = False + assert chain_monitor.watcher_asleep != watcher.asleep + + # Test that the Queue work + r_hash = get_random_value_hex(32) + chain_monitor.watcher_queue.put(r_hash) + assert watcher.block_queue.get() == r_hash + + +def test_attach_responder(chain_monitor): + responder = Responder(db_manager=None, chain_monitor=chain_monitor) + chain_monitor.attach_responder(responder.block_queue, responder.asleep) + + # Same kind of testing as with the attach watcher + assert responder.asleep == chain_monitor.watcher_asleep + responder.asleep = False + assert chain_monitor.watcher_asleep != responder.asleep + + r_hash = get_random_value_hex(32) + chain_monitor.responder_queue.put(r_hash) + assert responder.block_queue.get() == r_hash + + +def test_notify_subscribers(chain_monitor): + # Subscribers are only notified as long as they are awake + new_block = get_random_value_hex(32) + + # Queues should be empty to start with + assert chain_monitor.watcher_queue.empty() + assert chain_monitor.responder_queue.empty() + + chain_monitor.watcher_asleep = True + chain_monitor.responder_asleep = True + chain_monitor.notify_subscribers(new_block) + + # And remain empty afterwards since both subscribers where asleep + assert chain_monitor.watcher_queue.empty() + assert chain_monitor.responder_queue.empty() + + # Let's flag them as awake and try again + chain_monitor.watcher_asleep = False + chain_monitor.responder_asleep = False + chain_monitor.notify_subscribers(new_block) + + assert chain_monitor.watcher_queue.get() == new_block + assert chain_monitor.responder_queue.get() == new_block + + +def test_update_state(chain_monitor): + # The state is updated after receiving a new block (and only if the block is not already known). + # Let's start by setting a best_tip and a couple of old tips + new_block_hash = get_random_value_hex(32) + chain_monitor.best_tip = new_block_hash + chain_monitor.last_tips = [get_random_value_hex(32) for _ in range(5)] + + # Now we can try to update the state with an old best_tip and see how it doesn't work + assert chain_monitor.update_state(chain_monitor.last_tips[0]) is False + + # Same should happen with the current tip + assert chain_monitor.update_state(chain_monitor.best_tip) is False + + # The state should be correctly updated with a new block hash, the chain tip should change and the old tip should + # has been added to the last_tips + another_block_hash = get_random_value_hex(32) + assert chain_monitor.update_state(another_block_hash) is True + assert chain_monitor.best_tip == another_block_hash and new_block_hash == chain_monitor.last_tips[-1] + + +def test_monitor_chain_polling(): + # Try polling with the Watcher + chain_monitor = ChainMonitor() + chain_monitor.best_tip = BlockProcessor.get_best_block_hash() + + watcher = Watcher(db_manager=None, chain_monitor=chain_monitor, sk_der=None) + chain_monitor.attach_watcher(watcher.block_queue, asleep=False) + + # monitor_chain_polling runs until terminate if set + polling_thread = Thread(target=chain_monitor.monitor_chain_polling, kwargs={"polling_delta": 0.1}, daemon=True) + polling_thread.start() + + # Check that nothings changes as long as a block is not generated + for _ in range(5): + assert chain_monitor.watcher_queue.empty() + time.sleep(0.1) + + # And that it does if we generate a block + generate_block() + + chain_monitor.watcher_queue.get() + assert chain_monitor.watcher_queue.empty() + + chain_monitor.terminate = True + polling_thread.join() + + +def test_monitor_chain_zmq(): + # Try zmq with the Responder + chain_monitor = ChainMonitor() + chain_monitor.best_tip = BlockProcessor.get_best_block_hash() + + responder = Responder(db_manager=None, chain_monitor=chain_monitor) + chain_monitor.attach_responder(responder.block_queue, asleep=False) + + zmq_thread = Thread(target=chain_monitor.monitor_chain_zmq, daemon=True) + zmq_thread.start() + + # Queues should start empty + assert chain_monitor.responder_queue.empty() + + # And have a new block every time we generate one + for _ in range(3): + generate_block() + chain_monitor.responder_queue.get() + assert chain_monitor.responder_queue.empty() + + # If we flag it to sleep no notification is sent + chain_monitor.responder_asleep = True + + for _ in range(3): + generate_block() + assert chain_monitor.responder_queue.empty() + + chain_monitor.terminate = True + # The zmq thread needs a block generation to release from the recv method. + generate_block() + + zmq_thread.join() + + +def test_monitor_chain(): + # Not much to test here, this should launch two threads (one per monitor approach) and finish on terminate + chain_monitor = ChainMonitor() + + watcher = Watcher(db_manager=None, chain_monitor=chain_monitor, sk_der=None) + responder = Responder(db_manager=None, chain_monitor=chain_monitor) + chain_monitor.attach_responder(responder.block_queue, asleep=False) + chain_monitor.attach_watcher(watcher.block_queue, asleep=False) + + chain_monitor.best_tip = None + chain_monitor.monitor_chain() + + # The tip is updated before starting the threads, so it should have changed. + assert chain_monitor.best_tip is not None + + # Blocks should be received + for _ in range(5): + generate_block() + watcher_block = chain_monitor.watcher_queue.get() + responder_block = chain_monitor.responder_queue.get() + assert watcher_block == responder_block + assert chain_monitor.watcher_queue.empty() + assert chain_monitor.responder_queue.empty() + + # And the thread be terminated on terminate + chain_monitor.terminate = True + # The zmq thread needs a block generation to release from the recv method. + generate_block() From dfc90cd9306895605ba2de2b929752a2470e4fd8 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Tue, 7 Jan 2020 16:31:17 +0100 Subject: [PATCH 11/14] Sets ChainMonitor terminate on shutdown and removes unused imports --- pisa/pisad.py | 2 +- test/pisa/unit/test_api.py | 1 - 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/pisa/pisad.py b/pisa/pisad.py index b027662..ebf175a 100644 --- a/pisa/pisad.py +++ b/pisa/pisad.py @@ -8,7 +8,6 @@ from pisa.api import API from pisa.watcher import Watcher from pisa.builder import Builder from pisa.conf import BTC_NETWORK, PISA_SECRET_KEY -from pisa.responder import Responder from pisa.db_manager import DBManager from pisa.chain_monitor import ChainMonitor from pisa.block_processor import BlockProcessor @@ -20,6 +19,7 @@ logger = Logger("Daemon") def handle_signals(signal_received, frame): logger.info("Closing connection with appointments db") db_manager.db.close() + chain_monitor.terminate = True logger.info("Shutting down PISA") exit(0) diff --git a/test/pisa/unit/test_api.py b/test/pisa/unit/test_api.py index 1e1d4ef..fb236a0 100644 --- a/test/pisa/unit/test_api.py +++ b/test/pisa/unit/test_api.py @@ -10,7 +10,6 @@ from pisa.watcher import Watcher from pisa.tools import bitcoin_cli from pisa import HOST, PORT from pisa.conf import MAX_APPOINTMENTS -from pisa.chain_monitor import ChainMonitor from test.pisa.unit.conftest import ( generate_block, From 1a26d7d6a3418f3d9d79c5e46beeee45fd6b1a2c Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Mon, 13 Jan 2020 15:31:54 +0100 Subject: [PATCH 12/14] Fixes typos based on @orbitalturtle comments --- pisa/chain_monitor.py | 8 ++++---- test/pisa/unit/test_chain_monitor.py | 6 +++--- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/pisa/chain_monitor.py b/pisa/chain_monitor.py index 087096b..6327692 100644 --- a/pisa/chain_monitor.py +++ b/pisa/chain_monitor.py @@ -20,9 +20,9 @@ class ChainMonitor: Attributes: best_tip (:obj:`str`): a block hash representing the current best tip. - last_tips (:obj:`list`): a list of last chain tips. Used as an sliding window to avoid notifying about old tips. + last_tips (:obj:`list`): a list of last chain tips. Used as a sliding window to avoid notifying about old tips. terminate (:obj:`bool`): a flag to signal the termination of the :class:`ChainMonitor` (shutdown the tower). - check_tip (:obj:`Event`): an event that it's triggered at fixed time intervals and controls the polling thread. + check_tip (:obj:`Event`): an event that's triggered at fixed time intervals and controls the polling thread. lock (:obj:`Condition`): a lock used to protect concurrent access to the queues and ``best_tip`` by the zmq and polling threads. zmqSubSocket (:obj:`socket`): a socket to connect to ``bitcoind`` via ``zmq``. @@ -130,7 +130,7 @@ class ChainMonitor: while not self.terminate: self.check_tip.wait(timeout=polling_delta) - # Terminate could have been set wile the thread was blocked in wait + # Terminate could have been set while the thread was blocked in wait if not self.terminate: current_tip = BlockProcessor.get_best_block_hash() @@ -150,7 +150,7 @@ class ChainMonitor: while not self.terminate: msg = self.zmqSubSocket.recv_multipart() - # Terminate could have been set wile the thread was blocked in recv + # Terminate could have been set while the thread was blocked in recv if not self.terminate: topic = msg[0] body = msg[1] diff --git a/test/pisa/unit/test_chain_monitor.py b/test/pisa/unit/test_chain_monitor.py index d3037d8..4e8cdfb 100644 --- a/test/pisa/unit/test_chain_monitor.py +++ b/test/pisa/unit/test_chain_monitor.py @@ -13,7 +13,7 @@ from test.pisa.unit.conftest import get_random_value_hex, generate_block def test_init(run_bitcoind): # run_bitcoind is started here instead of later on to avoid race conditions while it initializes - # Not much to test here, just sanity checks to make sure nothings goes south in the future + # Not much to test here, just sanity checks to make sure nothing goes south in the future chain_monitor = ChainMonitor() assert chain_monitor.best_tip is None @@ -70,7 +70,7 @@ def test_notify_subscribers(chain_monitor): chain_monitor.responder_asleep = True chain_monitor.notify_subscribers(new_block) - # And remain empty afterwards since both subscribers where asleep + # And remain empty afterwards since both subscribers were asleep assert chain_monitor.watcher_queue.empty() assert chain_monitor.responder_queue.empty() @@ -115,7 +115,7 @@ def test_monitor_chain_polling(): polling_thread = Thread(target=chain_monitor.monitor_chain_polling, kwargs={"polling_delta": 0.1}, daemon=True) polling_thread.start() - # Check that nothings changes as long as a block is not generated + # Check that nothing changes as long as a block is not generated for _ in range(5): assert chain_monitor.watcher_queue.empty() time.sleep(0.1) From eb71ab8fc429a1c581834621bd249a04daf61135 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Mon, 13 Jan 2020 15:42:27 +0100 Subject: [PATCH 13/14] Adds missing args on docs and adds polling_delta parm to monitor_chain --- pisa/chain_monitor.py | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/pisa/chain_monitor.py b/pisa/chain_monitor.py index 6327692..689a223 100644 --- a/pisa/chain_monitor.py +++ b/pisa/chain_monitor.py @@ -125,6 +125,9 @@ class ChainMonitor: Monitors ``bitcoind`` via polling. Once the method is fired, it keeps monitoring as long as ``terminate`` is not set. Polling is performed once every ``polling_delta`` seconds. If a new best tip if found, the shared lock is acquired, the state is updated and the subscribers are notified, and finally the lock is released. + + Args: + polling_delta (:obj:`int`): the time delta between polls. """ while not self.terminate: @@ -164,13 +167,16 @@ class ChainMonitor: logger.info("New block received via zmq", block_hash=block_hash) self.lock.release() - def monitor_chain(self): + def monitor_chain(self, polling_delta=POLLING_DELTA): """ Main :class:`ChainMonitor` method. It initializes the ``best_tip`` to the current one (by querying the :obj:`BlockProcessor `) and creates two threads, one per each monitoring - approach (``zmq`` and ``polling``) + approach (``zmq`` and ``polling``). + + Args: + polling_delta (:obj:`int`): the time delta between polls by the ``monitor_chain_polling`` thread. """ self.best_tip = BlockProcessor.get_best_block_hash() - Thread(target=self.monitor_chain_polling, daemon=True).start() + Thread(target=self.monitor_chain_polling, daemon=True, kwargs={"polling_delta": polling_delta}).start() Thread(target=self.monitor_chain_zmq, daemon=True).start() From ae772bf91b05bf1173de88f54248d042bf343242 Mon Sep 17 00:00:00 2001 From: Sergi Delgado Segura Date: Mon, 13 Jan 2020 15:48:52 +0100 Subject: [PATCH 14/14] Adds missing tests Tests that given a block hash and the two monitor threads running, the hash is only notified to the subscribers once (by the first thread that notices it) --- test/pisa/unit/test_chain_monitor.py | 33 +++++++++++++++++++++++++++- 1 file changed, 32 insertions(+), 1 deletion(-) diff --git a/test/pisa/unit/test_chain_monitor.py b/test/pisa/unit/test_chain_monitor.py index 4e8cdfb..4c8b38c 100644 --- a/test/pisa/unit/test_chain_monitor.py +++ b/test/pisa/unit/test_chain_monitor.py @@ -97,7 +97,7 @@ def test_update_state(chain_monitor): assert chain_monitor.update_state(chain_monitor.best_tip) is False # The state should be correctly updated with a new block hash, the chain tip should change and the old tip should - # has been added to the last_tips + # have been added to the last_tips another_block_hash = get_random_value_hex(32) assert chain_monitor.update_state(another_block_hash) is True assert chain_monitor.best_tip == another_block_hash and new_block_hash == chain_monitor.last_tips[-1] @@ -192,3 +192,34 @@ def test_monitor_chain(): chain_monitor.terminate = True # The zmq thread needs a block generation to release from the recv method. generate_block() + + +def test_monitor_chain_single_update(): + # This test tests that if both threads try to add the same block to the queue, only the first one will make it + chain_monitor = ChainMonitor() + + watcher = Watcher(db_manager=None, chain_monitor=chain_monitor, sk_der=None) + responder = Responder(db_manager=None, chain_monitor=chain_monitor) + chain_monitor.attach_responder(responder.block_queue, asleep=False) + chain_monitor.attach_watcher(watcher.block_queue, asleep=False) + + chain_monitor.best_tip = None + + # We will create a block and wait for the polling thread. Then check the queues to see that the block hash has only + # been added once. + chain_monitor.monitor_chain(polling_delta=2) + generate_block() + + watcher_block = chain_monitor.watcher_queue.get() + responder_block = chain_monitor.responder_queue.get() + assert watcher_block == responder_block + assert chain_monitor.watcher_queue.empty() + assert chain_monitor.responder_queue.empty() + + # The delta for polling is 2 secs, so let's wait and see + time.sleep(2) + assert chain_monitor.watcher_queue.empty() + assert chain_monitor.responder_queue.empty() + + # We can also force an update and see that it won't go through + assert chain_monitor.update_state(watcher_block) is False