diff --git a/squeaknode/config/config.py b/squeaknode/config/config.py index e7d119c8..92c0a9bc 100644 --- a/squeaknode/config/config.py +++ b/squeaknode/config/config.py @@ -39,6 +39,8 @@ DEFAULT_SYNC_TIMEOUT_S = 10 DEFAULT_SYNC_BLOCK_INTERVAL = 2016 DEFAULT_SENT_OFFER_RETENTION_S = 86400 DEFAULT_SUBSCRIBE_INVOICES_RETRY_S = 10 +DEFAULT_SQUEAK_RETENTION_S = 604800 +DEFAULT_SQUEAK_DELETION_INTERVAL_S = 10 @section('bitcoin') @@ -101,6 +103,10 @@ class CoreConfig(Config): cast=int, required=False, default=DEFAULT_SENT_OFFER_RETENTION_S) subscribe_invoices_retry_s = key( cast=int, required=False, default=DEFAULT_SUBSCRIBE_INVOICES_RETRY_S) + squeak_retention_s = key( + cast=int, required=False, default=DEFAULT_SQUEAK_RETENTION_S) + squeak_deletion_interval_s = key( + cast=int, required=False, default=DEFAULT_SQUEAK_DELETION_INTERVAL_S) @section('sync') diff --git a/squeaknode/db/squeak_db.py b/squeaknode/db/squeak_db.py index 37812df2..0a014b07 100644 --- a/squeaknode/db/squeak_db.py +++ b/squeaknode/db/squeak_db.py @@ -92,6 +92,10 @@ class SqueakDb: def squeak_has_no_secret_key(self): return self.squeaks.c.secret_key == None # noqa: E711 + def squeak_is_older_than_retention(self, interval_s): + return self.datetime_now - timedelta(seconds=interval_s) > \ + self.squeaks.c.created + @property def profile_has_private_key(self): return self.profiles.c.private_key != None # noqa: E711 @@ -108,9 +112,9 @@ class SqueakDb: def datetime_now(self): return datetime.now(timezone.utc) - def squeak_newer_than_interval_s(self, interval_s): - return self.squeaks.c.created > \ - self.datetime_now - timedelta(seconds=interval_s) + # def squeak_newer_than_interval_s(self, interval_s): + # return self.squeaks.c.created > \ + # self.datetime_now - timedelta(seconds=interval_s) def received_offer_should_be_deleted(self): expire_time = ( @@ -425,6 +429,29 @@ class SqueakDb: hashes = [bytes.fromhex(row["hash"]) for row in rows] return hashes + def get_old_squeaks_to_delete( + self, + interval_s: int, + ) -> List[SqueakEntryWithProfile]: + """ Get squeaks older than retention that meet the + criteria for deletion. + """ + s = ( + select([self.squeaks, self.profiles]) + .select_from( + self.squeaks.outerjoin( + self.profiles, + self.profiles.c.address == self.squeaks.c.author_address, + ) + ) + .where(self.squeak_is_older_than_retention(interval_s)) + .where(self.profile_has_no_private_key) + ) + with self.get_connection() as connection: + result = connection.execute(s) + rows = result.fetchall() + return [self._parse_squeak_entry_with_profile(row) for row in rows] + def insert_profile(self, squeak_profile: SqueakProfile) -> int: """ Insert a new squeak profile. """ ins = self.profiles.insert().values( diff --git a/squeaknode/node/squeak_controller.py b/squeaknode/node/squeak_controller.py index 648ff97e..382453b4 100644 --- a/squeaknode/node/squeak_controller.py +++ b/squeaknode/node/squeak_controller.py @@ -519,3 +519,17 @@ class SqueakController: def reprocess_received_payments(self) -> None: self.squeak_db.clear_received_payment_settle_indices() self.payment_processor.start_processing() + + def delete_old_squeaks(self): + squeaks_to_delete = self.squeak_db.get_old_squeaks_to_delete( + self.config.core.squeak_retention_s, + ) + for squeak_entry_with_profile in squeaks_to_delete: + squeak = squeak_entry_with_profile.squeak_entry.squeak + squeak_hash = get_hash(squeak) + self.squeak_db.delete_squeak( + squeak_hash, + ) + logger.info("Deleted squeak: {}".format( + squeak_hash.hex(), + )) diff --git a/squeaknode/node/squeak_deletion_worker.py b/squeaknode/node/squeak_deletion_worker.py new file mode 100644 index 00000000..94b24a35 --- /dev/null +++ b/squeaknode/node/squeak_deletion_worker.py @@ -0,0 +1,28 @@ +import logging +import threading + +from squeaknode.node.squeak_controller import SqueakController + +logger = logging.getLogger(__name__) + + +class SqueakDeletionWorker: + def __init__( + self, + squeak_controller: SqueakController, + clean_interval_s: int, + ): + self.squeak_controller = squeak_controller + self.clean_interval_s = clean_interval_s + + def start_running(self): + timer = threading.Timer( + self.clean_interval_s, + self.start_running, + ) + timer.daemon = True + timer.start() + self.delete_old_squeaks() + + def delete_old_squeaks(self): + self.squeak_controller.delete_old_squeaks() diff --git a/squeaknode/node/squeak_node.py b/squeaknode/node/squeak_node.py index c0f374e6..61be1997 100644 --- a/squeaknode/node/squeak_node.py +++ b/squeaknode/node/squeak_node.py @@ -17,6 +17,7 @@ from squeaknode.lightning.lnd_lightning_client import LNDLightningClient from squeaknode.node.payment_processor import PaymentProcessor from squeaknode.node.process_received_payments_worker import ProcessReceivedPaymentsWorker from squeaknode.node.squeak_controller import SqueakController +from squeaknode.node.squeak_deletion_worker import SqueakDeletionWorker from squeaknode.node.squeak_offer_expiry_worker import SqueakOfferExpiryWorker from squeaknode.node.squeak_peer_sync_worker import SqueakPeerSyncWorker from squeaknode.node.squeak_rate_limiter import SqueakRateLimiter @@ -95,6 +96,10 @@ class SqueakNode: self.sent_offers_worker = ProcessReceivedPaymentsWorker( payment_processor, self.stopped, ) + self.squeak_deletion_worker = SqueakDeletionWorker( + squeak_controller, + self.config.core.squeak_deletion_interval_s, + ) handler = load_handler(squeak_controller) self.server = load_rpc_server( @@ -115,6 +120,7 @@ class SqueakNode: self.squeak_offer_expiry_worker.start_running() self.sent_offers_worker.start_running() + self.squeak_deletion_worker.start_running() # start peer rpc server if self.config.server.rpc_enabled: