From 3ee17aa3671cd46d3cae48acc5784bb7f838826a Mon Sep 17 00:00:00 2001 From: Jonathan Zernik Date: Sat, 1 Aug 2020 22:17:28 -0700 Subject: [PATCH] Run subscriber periodically (#203) --- squeakserver/node/squeak_node.py | 6 ++++++ .../node/squeak_subscription_downloader.py | 17 +++++++++++++++-- 2 files changed, 21 insertions(+), 2 deletions(-) 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()