diff --git a/squeakserver/node/squeak_node.py b/squeakserver/node/squeak_node.py index 55aa057d..123a73f5 100644 --- a/squeakserver/node/squeak_node.py +++ b/squeakserver/node/squeak_node.py @@ -12,6 +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_subscription_downloader import SqueakSubscriptionDownloader from squeakserver.server.buy_offer import BuyOffer from squeakserver.server.squeak_profile import SqueakProfile from squeakserver.server.squeak_subscription import SqueakSubscription @@ -53,10 +54,15 @@ class SqueakNode: self.squeak_rate_limiter, self.squeak_whitelist, ) + self.squeak_subscription_downloader = SqueakSubscriptionDownloader( + postgres_db, + self.squeak_store, + ) def start_running(self): # self.squeak_block_periodic_worker.start_running() self.squeak_block_queue_worker.start_running() + self.squeak_subscription_downloader.start_running() def save_uploaded_squeak(self, squeak): return self.squeak_store.save_uploaded_squeak(squeak) diff --git a/squeakserver/node/squeak_subscription_downloader.py b/squeakserver/node/squeak_subscription_downloader.py index ec6ea51e..ca6bb243 100644 --- a/squeakserver/node/squeak_subscription_downloader.py +++ b/squeakserver/node/squeak_subscription_downloader.py @@ -1,9 +1,22 @@ import logging import queue +import threading logger = logging.getLogger(__name__) +SUBSCRIBE_UPDATE_INTERVAL_S = 10.0 + + class SqueakSubscriptionDownloader: - def __init__(self): - pass + def __init__(self, postgres_db, squeak_store, update_interval_s=SUBSCRIBE_UPDATE_INTERVAL_S): + self.postgres_db = postgres_db + self.squeak_store = squeak_store + self.update_interval_s = update_interval_s + + def sync_subscriptions(self): + logger.info("Syncing subscriptions...") + + def start_running(self): + threading.Timer(self.update_interval_s, self.start_running).start() + self.sync_subscriptions()