From ccc812f997be53a823b0120a03bd12e559ff4233 Mon Sep 17 00:00:00 2001 From: Lee Chagolla-Christensen Date: Tue, 30 Sep 2025 22:54:50 -0700 Subject: [PATCH] simplify deletes, keep task alive --- libs/chromium/chromium/cookies.py | 16 ++-- .../file_enrichment/workflow_manager.py | 2 +- projects/housekeeping/housekeeping/main.py | 84 ++++++------------- 3 files changed, 34 insertions(+), 68 deletions(-) diff --git a/libs/chromium/chromium/cookies.py b/libs/chromium/chromium/cookies.py index f474ebe..ef78550 100644 --- a/libs/chromium/chromium/cookies.py +++ b/libs/chromium/chromium/cookies.py @@ -1,7 +1,7 @@ """Chromium Cookies file parsing and database operations.""" -import sqlite3 import asyncio +import sqlite3 import psycopg import structlog @@ -126,7 +126,7 @@ def _insert_cookies( try: value_dec_bytes = asyncio.run(dpapi_manager.decrypt_blob(Blob.parse(encrypted_value))) if value_dec_bytes: - value_dec = value_dec_bytes.decode('utf-8', errors='replace') + value_dec = value_dec_bytes.decode("utf-8", errors="replace") is_decrypted = True except: pass @@ -153,13 +153,15 @@ def _insert_cookies( # Or just 16-byte suffix value_dec_bytes = value_dec_bytes[:-16] - value_dec = value_dec_bytes.decode('utf-8', errors='replace') + value_dec = value_dec_bytes.decode("utf-8", errors="replace") is_decrypted = True except Exception as e: - logger.warning("Failed to decrypt cookie with state key", - state_key_id=state_key_id, - encryption_type=encryption_type, - error=str(e)) + logger.warning( + "Failed to decrypt cookie with state key", + state_key_id=state_key_id, + encryption_type=encryption_type, + error=str(e), + ) cookie_data = { "originating_object_id": file_enriched.object_id, diff --git a/projects/file_enrichment/file_enrichment/workflow_manager.py b/projects/file_enrichment/file_enrichment/workflow_manager.py index a238348..364aaba 100644 --- a/projects/file_enrichment/file_enrichment/workflow_manager.py +++ b/projects/file_enrichment/file_enrichment/workflow_manager.py @@ -299,7 +299,7 @@ class WorkflowManager: except Exception as e: logger.exception(e, message="Error resetting workflows in database", pid=os.getpid()) - logger.warning("WorkflowManager reset", active_count=len(self.active_workflows), pid=os.getpid()) + logger.info("WorkflowManager reset", active_count=len(self.active_workflows), pid=os.getpid()) return { "status": "success", diff --git a/projects/housekeeping/housekeeping/main.py b/projects/housekeeping/housekeeping/main.py index 1ae8ded..6e6382d 100644 --- a/projects/housekeeping/housekeeping/main.py +++ b/projects/housekeeping/housekeeping/main.py @@ -20,6 +20,7 @@ logger = structlog.get_logger(module=__name__) scheduler = AsyncIOScheduler() is_initialized = False storage = StorageMinio() +background_tasks = set() with DaprClient() as client: secret = client.get_secret(store_name="nemesis-secret-store", key="POSTGRES_CONNECTION_STRING") @@ -49,62 +50,24 @@ async def get_expired_object_ids(expiration_date: Optional[datetime] = None) -> """ try: conn = get_db_connection() - - # Use provided expiration date or current datetime comparison_date = expiration_date if expiration_date is not None else datetime.now() - # Define the WHERE clause based on the expiration_date - where_clause = "WHERE 1=1" # Default to get all records if datetime.max - params = tuple() - - if expiration_date != datetime.max: - where_clause = "WHERE expiration < %s" - params = (comparison_date,) - - # Get expired object_ids from the files table with conn.cursor() as cur: cur.execute( - f""" - SELECT object_id - FROM files - {where_clause} + """ + SELECT DISTINCT object_id + FROM ( + SELECT object_id FROM files WHERE expiration < %s + UNION + SELECT object_id FROM files_enriched WHERE expiration < %s + UNION + SELECT object_id FROM files_enriched_dataset WHERE expiration < %s + ) AS expired """, - params, + (comparison_date, comparison_date, comparison_date), ) - expired_files = cur.fetchall() + return [str(record[0]) for record in cur.fetchall()] - # Get expired object_ids from the files_enriched table - cur.execute( - f""" - SELECT object_id - FROM files_enriched - {where_clause} - """, - params, - ) - expired_files_enriched = cur.fetchall() - - # Get expired object_ids from the files_enriched_dataset table - cur.execute( - f""" - SELECT object_id - FROM files_enriched_dataset - {where_clause} - """, - params, - ) - expired_files_dataset = cur.fetchall() - - # Combine all object_ids - all_expired_object_ids = set() - for record in expired_files: - all_expired_object_ids.add(str(record[0])) - for record in expired_files_enriched: - all_expired_object_ids.add(str(record[0])) - for record in expired_files_dataset: - all_expired_object_ids.add(str(record[0])) - - return list(all_expired_object_ids) except Exception as e: logger.exception(e, message="Error getting expired object IDs from database") return [] @@ -252,7 +215,7 @@ async def delete_expired_containers(expiration_date: Optional[datetime] = None) logger.info( "Successfully deleted expired containers", container_count=count_to_delete, - expiration_date=comparison_date if expiration_date != datetime.max else "all" + expiration_date=comparison_date if expiration_date != datetime.max else "all", ) return True @@ -287,42 +250,41 @@ async def run_cleanup_job(expiration_date: Optional[datetime] = None): for x in range(3): # Step 1: Get expired object IDs from database using the provided expiration date expired_object_ids = await get_expired_object_ids(expiration_date) - logger.info("Found expired objects", count=len(expired_object_ids), round=(x+1)) + logger.info("Found expired objects", count=len(expired_object_ids), round=(x + 1)) if not expired_object_ids: - logger.info("No expired objects found, cleanup job completed", round=(x+1)) + logger.info("No expired objects found, cleanup job completed", round=(x + 1)) return # Step 2: Get related transform object IDs transform_object_ids = await get_transform_object_ids(expired_object_ids) - logger.info("Found related transform objects", count=len(transform_object_ids), round=(x+1)) + logger.info("Found related transform objects", count=len(transform_object_ids), round=(x + 1)) # Step 3: Combine all object IDs that need to be deleted from Minio all_object_ids = list(set(expired_object_ids + transform_object_ids)) # Step 4: Delete objects from Minio using the StorageMinio instance deleted_count = storage.delete_objects(all_object_ids) - logger.info("Deleted objects from Minio", count=deleted_count, total=len(all_object_ids), round=(x+1)) + logger.info("Deleted objects from Minio", count=deleted_count, total=len(all_object_ids), round=(x + 1)) # Step 5: Delete database entries db_delete_success = await delete_database_entries(expired_object_ids) if db_delete_success: - logger.info("Successfully deleted database entries", round=(x+1)) + logger.info("Successfully deleted database entries", round=(x + 1)) else: - logger.error("Failed to delete some database entries during cleanup job", round=(x+1)) + logger.error("Failed to delete some database entries during cleanup job", round=(x + 1)) # Step 6: delete any expired containers container_delete_success = await delete_expired_containers(expiration_date) if container_delete_success: - logger.info("Successfully deleted container entries", round=(x+1)) + logger.info("Successfully deleted container entries", round=(x + 1)) else: - logger.error("Failed to delete some container entries during cleanup job", round=(x+1)) + logger.error("Failed to delete some container entries during cleanup job", round=(x + 1)) await asyncio.sleep(20) logger.info("Cleanup job complete") - except Exception as e: logger.exception(e, message="Error running cleanup job") @@ -436,7 +398,9 @@ async def trigger_cleanup(request: CleanupRequest): } # Trigger the cleanup job with the specified expiration - asyncio.create_task(run_cleanup_job(expiration_date)) + task = asyncio.create_task(run_cleanup_job(expiration_date)) + background_tasks.add(task) + task.add_done_callback(background_tasks.discard) return { "message": "Cleanup job triggered successfully",