Get offers in sync thread (#240)

* Get offers in sync background thread

* Fix get squeaks needing offer database query

* Fix getting squeak download and offer in signle sync

* Remove get offer controller

* Remove old comments and log lines

* Remove load offers rpc method
This commit is contained in:
Jonathan Zernik 2020-08-15 17:41:09 -07:00 committed by GitHub
parent 1394c0b620
commit 86a3bd5ba1
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
10 changed files with 112 additions and 128 deletions

View file

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

View file

@ -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) {}

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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,
)

View file

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

View file

@ -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 = """