diff --git a/itests/tests/test_squeak_node.py b/itests/tests/test_squeak_node.py index e64323a4..70009bfe 100644 --- a/itests/tests/test_squeak_node.py +++ b/itests/tests/test_squeak_node.py @@ -557,12 +557,3 @@ def test_delete_peer(server_stub, admin_stub, peer_id): ) ) assert "Peer not found." in str(excinfo.value) - - -def test_load_buy_offers(server_stub, admin_stub, saved_squeak_hash): - # Load buy offers for the squeak - admin_stub.LoadBuyOffers( - squeak_admin_pb2.LoadBuyOffersRequest( - squeak_hash=saved_squeak_hash.hex(), - ) - ) diff --git a/proto/squeak_admin.proto b/proto/squeak_admin.proto index f09fde41..58180057 100644 --- a/proto/squeak_admin.proto +++ b/proto/squeak_admin.proto @@ -108,10 +108,6 @@ service SqueakAdmin { */ rpc DeletePeer (DeletePeerRequest) returns (DeletePeerReply) {} - /** sqkadmin: `loadbuyoffers` - */ - rpc LoadBuyOffers (LoadBuyOffersRequest) returns (LoadBuyOffersReply) {} - /** sqkadmin: `getbuyoffers` */ rpc GetBuyOffers (GetBuyOffersRequest) returns (GetBuyOffersReply) {} diff --git a/squeakserver/admin/squeak_admin_server_handler.py b/squeakserver/admin/squeak_admin_server_handler.py index 8881094d..e1fd1fa0 100644 --- a/squeakserver/admin/squeak_admin_server_handler.py +++ b/squeakserver/admin/squeak_admin_server_handler.py @@ -200,10 +200,6 @@ class SqueakAdminServerHandler(object): logger.info("Handle delete squeak peer with id: {}".format(peer_id)) self.squeak_node.delete_peer(peer_id) - def handle_load_buy_offers(self, squeak_hash_str): - logger.info("Handle load buy offers for hash: {}".format(squeak_hash_str)) - self.squeak_node.load_buy_offers(squeak_hash_str) - def handle_get_buy_offers(self, squeak_hash_str): logger.info("Handle get buy offers for hash: {}".format(squeak_hash_str)) return self.squeak_node.get_buy_offers_with_peer(squeak_hash_str) diff --git a/squeakserver/admin/squeak_admin_server_servicer.py b/squeakserver/admin/squeak_admin_server_servicer.py index d97d1bb6..65640df0 100644 --- a/squeakserver/admin/squeak_admin_server_servicer.py +++ b/squeakserver/admin/squeak_admin_server_servicer.py @@ -218,11 +218,6 @@ class SqueakAdminServerServicer(squeak_admin_pb2_grpc.SqueakAdminServicer): self.handler.handle_delete_squeak_peer(peer_id) return squeak_admin_pb2.DeletePeerReply() - def LoadBuyOffers(self, request, context): - squeak_hash_str = request.squeak_hash - self.handler.handle_load_buy_offers(squeak_hash_str) - return squeak_admin_pb2.LoadBuyOffersReply() - def GetBuyOffers(self, request, context): squeak_hash_str = request.squeak_hash offers = self.handler.handle_get_buy_offers(squeak_hash_str) diff --git a/squeakserver/node/peer_download.py b/squeakserver/node/peer_download.py index 4697d211..597086ed 100644 --- a/squeakserver/node/peer_download.py +++ b/squeakserver/node/peer_download.py @@ -3,6 +3,7 @@ import threading from squeakserver.server.util import get_hash from squeakserver.network.peer_client import PeerClient +from squeakserver.node.peer_get_offer import PeerGetOffer logger = logging.getLogger(__name__) @@ -17,12 +18,14 @@ class PeerDownload: block_height, squeak_store, postgres_db, + lightning_client, lookup_block_interval=LOOKUP_BLOCK_INTERVAL, ): self.peer = peer self.block_height = block_height self.squeak_store = squeak_store self.postgres_db = postgres_db + self.lightning_client = lightning_client self.lookup_block_interval = lookup_block_interval self.peer_client = PeerClient( @@ -51,7 +54,7 @@ class PeerDownload: if self.stopped(): return - # Get local hashes + # Get local hashes of downloaded squeaks local_hashes = self._get_local_hashes(addresses, min_block, max_block) logger.debug("Got local hashes: {}".format(len(local_hashes))) for hash in local_hashes: @@ -73,6 +76,28 @@ class PeerDownload: return self._download_squeak(hash) + if self.stopped(): + return + + # Get local hashes of locked squeaks that don't have an offer from this peer. + locked_hashes = self._get_locked_hashes(addresses, min_block, max_block) + logger.debug("Got locked hashes: {}".format(len(locked_hashes))) + for hash in locked_hashes: + logger.debug("locked hash: {}".format(hash.hex())) + + # Get hashes to get offer + hashes_to_get_offer = set(remote_hashes) & set(locked_hashes) + logger.debug("Hashes to get offer: {}".format(len(hashes_to_get_offer))) + for hash in hashes_to_get_offer: + logger.debug("hash to get offer: {}".format(hash.hex())) + + # Download offers for the hashes + # TODO: catch exception downloading individual squeak + for hash in hashes_to_get_offer: + if self.stopped(): + return + self._download_offer(hash) + def stop(self): self._stop_event.set() @@ -80,7 +105,19 @@ class PeerDownload: return self._stop_event.is_set() def _get_local_hashes(self, addresses, min_block, max_block): - return self.squeak_store.lookup_squeaks(addresses, min_block, max_block) + return self.squeak_store.lookup_squeaks_include_locked( + addresses, + min_block, + max_block, + ) + + def _get_locked_hashes(self, addresses, min_block, max_block): + return self.squeak_store.lookup_squeaks_needing_offer( + addresses, + min_block, + max_block, + self.peer.peer_id, + ) def _get_remote_hashes(self, addresses, min_block, max_block): return self.peer_client.lookup_squeaks(addresses, min_block, max_block) @@ -99,3 +136,14 @@ class PeerDownload: profile.address for profile in followed_profiles ] + + def _download_offer(self, squeak_hash): + logger.info("Downloading offer for hash: {}".format(squeak_hash.hex())) + peer_get_offer = PeerGetOffer( + self.peer, + squeak_hash, + self.squeak_store, + self.postgres_db, + self.lightning_client, + ) + peer_get_offer.get_offer() diff --git a/squeakserver/node/squeak_get_offer_controller.py b/squeakserver/node/squeak_get_offer_controller.py deleted file mode 100644 index 8797c3f8..00000000 --- a/squeakserver/node/squeak_get_offer_controller.py +++ /dev/null @@ -1,86 +0,0 @@ -import logging -import threading - -from squeakserver.server.util import get_hash -from squeakserver.node.peer_download import PeerDownload -from squeakserver.node.peer_get_offer import PeerGetOffer - -logger = logging.getLogger(__name__) - - -HOUR_IN_SECONDS = 3600 - - -class SqueakGetOfferStatus: - def __init__(self): - self.downloads = {} - - def add_download(self, peer, squeak_hash, peer_download): - key = (peer.peer_id, squeak_hash) - self.downloads[key] = peer_download - - def is_downloading(self, peer, squeak_hash): - key = (peer.peer_id, squeak_hash) - return key in self.downloads - - def remove_download(self, peer, squeak_hash): - key = (peer.peer_id, squeak_hash) - del self.downloads[key] - - def get_current_downloads(self): - return self.downloads.values() - - -class SqueakGetOfferController: - def __init__(self, squeak_store, postgres_db, lightning_client): - self.squeak_get_offer_status = SqueakGetOfferStatus() - self.squeak_store = squeak_store - self.postgres_db = postgres_db - self.lightning_client = lightning_client - - def get_offers(self, peers, squeak_hash): - self._download_from_peers(peers, squeak_hash) - - def _download_from_peers(self, peers, squeak_hash): - for peer in peers: - if peer.downloading: - download_thread = threading.Thread( - target=self._download_from_peer, - args=(peer, squeak_hash,), - ) - download_thread.start() - - def _download_from_peer(self, peer, squeak_hash): - peer_get_offer = PeerGetOffer( - peer, - squeak_hash, - self.squeak_store, - self.postgres_db, - self.lightning_client, - ) - try: - logger.debug("Trying to get offer from peer: {} for squeak: {}".format(peer.peer_id, squeak_hash)) - with self.DownloadingContextManager(peer, squeak_hash, peer_get_offer, self.squeak_get_offer_status) as downloading_manager: - peer_get_offer.get_offer() - except Exception as e: - logger.error("Get offer from peer failed.", exc_info=True) - - class DownloadingContextManager(): - def __init__(self, peer, squeak_hash, peer_download, squeak_get_offer_status): - self.peer = peer - self.squeak_hash = squeak_hash - self.peer_download = peer_download - self.squeak_get_offer_status = squeak_get_offer_status - - if self.squeak_get_offer_status.is_downloading(self.peer, self.squeak_hash): - raise Exception("Peer {} is already getting offer for squeak hash: {}".format( - self.peer.peer_id, - self.squeak_hash, - )) - - def __enter__(self): - self.squeak_get_offer_status.add_download(self.peer, self.squeak_hash, self.peer_download) - return self - - def __exit__(self, exc_type, exc_value, exc_traceback): - self.squeak_get_offer_status.remove_download(self.peer, self.squeak_hash) diff --git a/squeakserver/node/squeak_node.py b/squeakserver/node/squeak_node.py index f8fb96fa..d7cffd67 100644 --- a/squeakserver/node/squeak_node.py +++ b/squeakserver/node/squeak_node.py @@ -14,7 +14,6 @@ from squeakserver.node.squeak_whitelist import SqueakWhitelist from squeakserver.node.squeak_store import SqueakStore from squeakserver.node.squeak_peer_sync_worker import SqueakPeerSyncWorker from squeakserver.node.squeak_sync_status import SqueakSyncController -from squeakserver.node.squeak_get_offer_controller import SqueakGetOfferController from squeakserver.server.buy_offer import BuyOffer from squeakserver.server.squeak_profile import SqueakProfile from squeakserver.server.squeak_peer import SqueakPeer @@ -60,16 +59,12 @@ class SqueakNode: self.blockchain_client, self.squeak_store, self.postgres_db, + self.lightning_client, ) self.squeak_peer_sync_worker = SqueakPeerSyncWorker( postgres_db, self.squeak_sync_controller, ) - self.squeak_get_offer_controller = SqueakGetOfferController( - self.squeak_store, - self.postgres_db, - self.lightning_client, - ) def start_running(self): self.squeak_block_periodic_worker.start_running() @@ -241,11 +236,6 @@ class SqueakNode: def delete_peer(self, peer_id): self.postgres_db.delete_peer(peer_id) - def load_buy_offers(self, squeak_hash_str): - peers = self.postgres_db.get_peers() - squeak_hash = bytes.fromhex(squeak_hash_str) - self.squeak_get_offer_controller.get_offers(peers, squeak_hash) - def get_buy_offers_with_peer(self, squeak_hash_str): squeak_hash = bytes.fromhex(squeak_hash_str) return self.postgres_db.get_offers_with_peer(squeak_hash_str) diff --git a/squeakserver/node/squeak_store.py b/squeakserver/node/squeak_store.py index 7610f992..0776dbfb 100644 --- a/squeakserver/node/squeak_store.py +++ b/squeakserver/node/squeak_store.py @@ -26,11 +26,14 @@ class SqueakStore: raise Exception("Excedeed allowed number of squeaks per block.") inserted_squeak_hash = self.postgres_db.insert_squeak(squeak) - self.squeak_block_verifier.add_squeak_to_queue(inserted_squeak_hash) + # self.squeak_block_verifier.add_squeak_to_queue(inserted_squeak_hash) + # Slow operation because of blockchain lookup + self.squeak_block_verifier.verify_squeak_block(inserted_squeak_hash) return inserted_squeak_hash def save_created_squeak(self, squeak): inserted_squeak_hash = self.postgres_db.insert_squeak(squeak) + # Slow operation because of blockchain lookup self.squeak_block_verifier.verify_squeak_block(inserted_squeak_hash) return inserted_squeak_hash @@ -68,4 +71,24 @@ class SqueakStore: return self.postgres_db.delete_squeak(squeak_hash) def lookup_squeaks(self, addresses, min_block, max_block): - return self.postgres_db.lookup_squeaks(addresses, min_block, max_block) + return self.postgres_db.lookup_squeaks( + addresses, + min_block, + max_block, + ) + + def lookup_squeaks_include_locked(self, addresses, min_block, max_block): + return self.postgres_db.lookup_squeaks( + addresses, + min_block, + max_block, + include_locked=True, + ) + + def lookup_squeaks_needing_offer(self, addresses, min_block, max_block, peer_id): + return self.postgres_db.lookup_squeaks_needing_offer( + addresses, + min_block, + max_block, + peer_id, + ) diff --git a/squeakserver/node/squeak_sync_status.py b/squeakserver/node/squeak_sync_status.py index 76205347..3948dd5c 100644 --- a/squeakserver/node/squeak_sync_status.py +++ b/squeakserver/node/squeak_sync_status.py @@ -42,11 +42,12 @@ class SqueakSyncStatus: class SqueakSyncController: - def __init__(self, blockchain_client, squeak_store, postgres_db): + def __init__(self, blockchain_client, squeak_store, postgres_db, lightning_client): self.squeak_sync_status = SqueakSyncStatus() self.blockchain_client = blockchain_client self.squeak_store = squeak_store self.postgres_db = postgres_db + self.lightning_client = lightning_client def sync_peers(self, peers): try: @@ -82,6 +83,7 @@ class SqueakSyncController: block_height, self.squeak_store, self.postgres_db, + self.lightning_client, ) try: logger.debug("Trying to download from peer: {}".format(peer)) diff --git a/squeakserver/server/postgres_db.py b/squeakserver/server/postgres_db.py index a6ba3918..e94fdd67 100644 --- a/squeakserver/server/postgres_db.py +++ b/squeakserver/server/postgres_db.py @@ -78,7 +78,7 @@ class PostgresDb: squeak.encContent.hex(), squeak.vchScriptSig, str(squeak.GetAddress()), - squeak.vchDecryptionKey, + squeak.GetDecryptionKey().get_bytes() if squeak.HasDecryptionKey() else None, ), ) # get the generated hash back @@ -169,14 +169,14 @@ class PostgresDb: rows = curs.fetchall() return [self._parse_squeak_entry_with_profile(row) for row in rows] - def lookup_squeaks(self, addresses, min_block, max_block, include_unverified=False): + def lookup_squeaks(self, addresses, min_block, max_block, include_unverified=False, include_locked=False): """ Lookup squeaks. """ sql = """ SELECT hash FROM squeak WHERE author_address IN %s AND n_block_height >= %s AND n_block_height <= %s - AND vch_decryption_key IS NOT NULL + AND (vch_decryption_key IS NOT NULL) OR %s AND ((block_header IS NOT NULL) OR %s); """ addresses_tuple = tuple(addresses) @@ -188,7 +188,7 @@ class PostgresDb: # mogrify to debug. # logger.info(curs.mogrify(sql, (addresses_tuple, min_block, max_block))) curs.execute( - sql, (addresses_tuple, min_block, max_block, include_unverified) + sql, (addresses_tuple, min_block, max_block, include_locked, include_unverified) ) rows = curs.fetchall() hashes = [bytes.fromhex(row["hash"]) for row in rows] @@ -218,6 +218,35 @@ class PostgresDb: hashes = [bytes.fromhex(row["hash"]) for row in rows] return hashes + def lookup_squeaks_needing_offer(self, addresses, min_block, max_block, peer_id, include_unverified=False): + """ Lookup squeaks that are locked and don't have an offer. """ + sql = """ + SELECT hash FROM squeak + LEFT JOIN offer + ON squeak.hash=offer.squeak_hash + AND offer.peer_id=%s + WHERE author_address IN %s + AND n_block_height >= %s + AND n_block_height <= %s + AND vch_decryption_key IS NULL + AND ((block_header IS NOT NULL) OR %s) + AND offer.squeak_hash IS NULL + """ + addresses_tuple = tuple(addresses) + + if not addresses: + return [] + + with self.get_cursor() as curs: + # mogrify to debug. + # logger.info(curs.mogrify(sql, (addresses_tuple, min_block, max_block))) + curs.execute( + sql, (peer_id, addresses_tuple, min_block, max_block, include_unverified) + ) + rows = curs.fetchall() + hashes = [bytes.fromhex(row["hash"]) for row in rows] + return hashes + def insert_profile(self, squeak_profile): """ Insert a new squeak profile. """ sql = """