squeaknode/squeakserver/node/squeak_sync_status.py
Jonathan Zernik 86a3bd5ba1
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
2020-08-15 17:41:09 -07:00

139 lines
4.7 KiB
Python

import logging
import threading
from squeakserver.server.util import get_hash
from squeakserver.node.peer_download import PeerDownload
from squeakserver.node.peer_upload import PeerUpload
logger = logging.getLogger(__name__)
HOUR_IN_SECONDS = 3600
class SqueakSyncStatus:
def __init__(self):
self.downloads = {}
self.uploads = {}
def add_download(self, peer, peer_download):
self.downloads[peer.peer_id] = peer_download
def add_upload(self, peer, peer_upload):
self.uploads[peer.peer_id] = peer_upload
def is_downloading(self, peer):
return peer.peer_id in self.downloads
def is_uploading(self, peer):
return peer.peer_id in self.uploads
def remove_download(self, peer):
del self.downloads[peer.peer_id]
def remove_upload(self, peer):
del self.uploads[peer.peer_id]
def get_current_downloads(self):
return self.downloads.values()
def get_current_uploads(self):
return self.uploads.values()
class SqueakSyncController:
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:
block_info = self.blockchain_client.get_best_block_info()
block_height = block_info.block_height
except Exception as e:
logger.error("Failed to sync because unable to get blockchain info.", exc_info=True)
return
self._download_from_peers(peers, block_height)
self._upload_to_peers(peers, block_height)
def _download_from_peers(self, peers, block_height):
for peer in peers:
if peer.downloading:
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:
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(
peer,
block_height,
self.squeak_store,
self.postgres_db,
self.lightning_client,
)
try:
logger.debug("Trying to download from peer: {}".format(peer))
with self.DownloadingContextManager(peer, peer_download, self.squeak_sync_status) as downloading_manager:
peer_download.download()
except Exception as e:
logger.error("Download from peer failed.", exc_info=True)
def _upload_to_peer(self, peer, block_height):
peer_upload = PeerUpload(
peer,
block_height,
self.squeak_store,
self.postgres_db,
)
try:
logger.debug("Trying to upload to peer: {}".format(peer))
with self.UploadingContextManager(peer, peer_upload, self.squeak_sync_status) as uploading_manager:
peer_upload.upload()
except Exception as e:
logger.error("Upload from peer failed.", exc_info=True)
class DownloadingContextManager():
def __init__(self, peer, peer_download, squeak_sync_status):
self.peer = peer
self.peer_download = peer_download
self.squeak_sync_status = squeak_sync_status
if self.squeak_sync_status.is_downloading(self.peer):
raise Exception("Peer is already downloading: {}".format(self.peer))
def __enter__(self):
self.squeak_sync_status.add_download(self.peer, self.peer_download)
return self
def __exit__(self, exc_type, exc_value, exc_traceback):
self.squeak_sync_status.remove_download(self.peer)
class UploadingContextManager():
def __init__(self, peer, peer_upload, squeak_sync_status):
self.peer = peer
self.peer_upload = peer_upload
self.squeak_sync_status = squeak_sync_status
if self.squeak_sync_status.is_uploading(self.peer):
raise Exception("Peer is already uploading: {}".format(self.peer))
def __enter__(self):
self.squeak_sync_status.add_upload(self.peer, self.peer_upload)
return self
def __exit__(self, exc_type, exc_value, exc_traceback):
self.squeak_sync_status.remove_upload(self.peer)