Improve setup download upload threads (#228)

* Improve setup of download and upload threads

* Rename squeak peer sync worker class
This commit is contained in:
Jonathan Zernik 2020-08-05 02:44:22 -07:00 committed by GitHub
parent ba3b7f9b2d
commit fc08f74801
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
6 changed files with 19 additions and 11 deletions

View file

@ -15,7 +15,7 @@ logger = logging.getLogger(__name__)
# we need to use that cipher suite otherwise there will be a handhsake
# error when we communicate with the lnd rpc server.
os.environ["GRPC_SSL_CIPHER_SUITES"] = "HIGH+ECDSA"
os.environ["GRPC_VERBOSITY"] = "DEBUG"
# os.environ["GRPC_VERBOSITY"] = "DEBUG"
class LNDLightningClient:

View file

@ -82,6 +82,7 @@ class PeerDownload:
self.squeak_store.save_downloaded_squeak(squeak)
def _download_squeak(self, squeak_hash):
logger.info("Downloading squeak: {}".format(squeak_hash.hex()))
squeak = self.peer_client.get_squeak(squeak_hash)
self._save_squeak(squeak)

View file

@ -83,6 +83,7 @@ class PeerUpload:
return squeak_entry.squeak
def _upload_squeak(self, squeak_hash):
logger.info("Uploading squeak: {}".format(squeak_hash.hex()))
squeak = self._get_local_squeak(squeak_hash)
self.peer_client.post_squeak(squeak)

View file

@ -12,7 +12,7 @@ from squeakserver.node.squeak_maker import SqueakMaker
from squeakserver.node.squeak_rate_limiter import SqueakRateLimiter
from squeakserver.node.squeak_whitelist import SqueakWhitelist
from squeakserver.node.squeak_store import SqueakStore
from squeakserver.node.squeak_peer_downloader import SqueakPeerDownloader
from squeakserver.node.squeak_peer_sync_worker import SqueakPeerSyncWorker
from squeakserver.node.squeak_sync_status import SqueakSyncController
from squeakserver.server.buy_offer import BuyOffer
from squeakserver.server.squeak_profile import SqueakProfile
@ -60,7 +60,7 @@ class SqueakNode:
self.squeak_store,
self.postgres_db,
)
self.squeak_peer_downloader = SqueakPeerDownloader(
self.squeak_peer_sync_worker = SqueakPeerSyncWorker(
postgres_db,
self.squeak_sync_controller,
)
@ -68,7 +68,7 @@ class SqueakNode:
def start_running(self):
self.squeak_block_periodic_worker.start_running()
self.squeak_block_queue_worker.start_running()
self.squeak_peer_downloader.start_running()
self.squeak_peer_sync_worker.start_running()
def save_uploaded_squeak(self, squeak):
return self.squeak_store.save_uploaded_squeak(squeak)

View file

@ -8,7 +8,7 @@ logger = logging.getLogger(__name__)
SUBSCRIBE_UPDATE_INTERVAL_S = 10.0
class SqueakPeerDownloader:
class SqueakPeerSyncWorker:
def __init__(self,
postgres_db,
squeak_sync_controller,

View file

@ -61,12 +61,20 @@ class SqueakSyncController:
def _download_from_peers(self, peers, block_height):
for peer in peers:
if peer.downloading:
self._download_from_peer(peer, block_height)
download_thread = threading.Thread(
target=self._download_from_peer,
args=(peer, block_height,),
)
download_thread.start()
def _upload_to_peers(self, peers, block_height):
for peer in peers:
if peer.uploading:
self._upload_to_peer(peer, block_height)
upload_thread = threading.Thread(
target=self._upload_to_peer,
args=(peer, block_height,),
)
upload_thread.start()
def _download_from_peer(self, peer, block_height):
peer_download = PeerDownload(
@ -78,8 +86,7 @@ class SqueakSyncController:
try:
logger.debug("Trying to download from peer: {}".format(peer))
with self.DownloadingContextManager(peer, peer_download, self.squeak_sync_status) as downloading_manager:
download_thread = threading.Thread(target=peer_download.download)
download_thread.start()
peer_download.download()
except Exception as e:
logger.error("Download from peer failed.", exc_info=True)
@ -93,8 +100,7 @@ class SqueakSyncController:
try:
logger.debug("Trying to upload to peer: {}".format(peer))
with self.UploadingContextManager(peer, peer_upload, self.squeak_sync_status) as uploading_manager:
upload_thread = threading.Thread(target=peer_upload.upload)
upload_thread.start()
peer_upload.upload()
except Exception as e:
logger.error("Upload from peer failed.", exc_info=True)