From 6a3dd08728ac8d04fd4084fa3cc721789f6a5a0c Mon Sep 17 00:00:00 2001 From: Russell Ballestrini Date: Fri, 8 Aug 2025 21:40:13 -0400 Subject: [PATCH] Optimize scheduler for minimal SQLite database locks - Redesigned scheduler to use separate short-lived sessions for lock operations - Lock acquisition and release now use independent database connections - Each scheduled action processed in its own transaction to minimize lock time - Error handling isolates failed actions to prevent blocking others - Re-enabled scheduler with optimized locking strategy Key improvements: - acquire_scheduler_lock() uses own connection, disposes immediately - release_scheduler_lock() uses own connection, disposes immediately - process_single_scheduled_action() uses transaction-per-action approach - Prevents long-running transactions that were blocking web requests russell@unturf.com is the boss --- .gitlab-ci.yml | 2 +- app.py | 195 ++++++++++++++++++++++++++++++------------------- 2 files changed, 120 insertions(+), 77 deletions(-) diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index 1eeb511..4ee1da3 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -63,7 +63,7 @@ deploy: [Service] WorkingDirectory=${APP_DIR} - Environment="PATH=${VENV_DIR}/bin" "DISABLE_SCHEDULER=true" + Environment="PATH=${VENV_DIR}/bin" ExecStart=${VENV_DIR}/bin/uwsgi \\ --master \\ --enable-threads \\ diff --git a/app.py b/app.py index 137b968..588e564 100644 --- a/app.py +++ b/app.py @@ -622,49 +622,72 @@ def reader_required(view_func): ################################################################################ -def acquire_scheduler_lock(dbsession): - """Acquire the scheduler lock.""" +def acquire_scheduler_lock(): + """Acquire the scheduler lock with minimal transaction time.""" process_id = f"{os.getpid()}_{int(time.time())}" + + # Use a separate short-lived session just for the lock + main_engine = create_engine( + DB_URL, + connect_args={"check_same_thread": False}, + poolclass=StaticPool, + ) + MainSessionFactory = sessionmaker(bind=main_engine) + + try: + with MainSessionFactory() as dbsession: + # Try to get existing lock + lock = dbsession.query(SchedulerLock).first() - # Try to get existing lock - lock = dbsession.query(SchedulerLock).first() + if not lock: + # Create new lock + lock = SchedulerLock( + locked=True, locked_at=datetime.datetime.utcnow(), process_id=process_id + ) + dbsession.add(lock) + dbsession.commit() + return True - if not lock: - # Create new lock - lock = SchedulerLock( - locked=True, locked_at=datetime.datetime.utcnow(), process_id=process_id - ) - dbsession.add(lock) - dbsession.commit() - return True + # Check if lock is stale (older than 5 minutes) + if lock.locked and lock.locked_at: + if datetime.datetime.utcnow() - lock.locked_at > datetime.timedelta(minutes=5): + log.warning("Detected stale scheduler lock, releasing it") + lock.locked = False + lock.locked_at = None + lock.process_id = None - # Check if lock is stale (older than 5 minutes) - if lock.locked and lock.locked_at: - if datetime.datetime.utcnow() - lock.locked_at > datetime.timedelta(minutes=5): - log.warning("Detected stale scheduler lock, releasing it") - lock.locked = False - lock.locked_at = None - lock.process_id = None - dbsession.commit() + if not lock.locked: + lock.locked = True + lock.locked_at = datetime.datetime.utcnow() + lock.process_id = process_id + dbsession.commit() + return True - if not lock.locked: - lock.locked = True - lock.locked_at = datetime.datetime.utcnow() - lock.process_id = process_id - dbsession.commit() - return True - - return False + return False + finally: + main_engine.dispose() -def release_scheduler_lock(dbsession): - """Release the scheduler lock.""" - lock = dbsession.query(SchedulerLock).first() - if lock and lock.locked: - lock.locked = False - lock.locked_at = None - lock.process_id = None - dbsession.commit() +def release_scheduler_lock(): + """Release the scheduler lock with minimal transaction time.""" + # Use a separate short-lived session just for the lock + main_engine = create_engine( + DB_URL, + connect_args={"check_same_thread": False}, + poolclass=StaticPool, + ) + MainSessionFactory = sessionmaker(bind=main_engine) + + try: + with MainSessionFactory() as dbsession: + lock = dbsession.query(SchedulerLock).first() + if lock and lock.locked: + lock.locked = False + lock.locked_at = None + lock.process_id = None + dbsession.commit() + finally: + main_engine.dispose() def process_scheduled_actions(): @@ -678,23 +701,27 @@ def process_scheduled_actions(): main_dbsession = MainSessionFactory() try: - if not acquire_scheduler_lock(main_dbsession): + if not acquire_scheduler_lock(): return # Another process is already running log.info("Scheduler acquired lock, processing scheduled actions...") - # Get all namespaces + # Get all namespaces with the existing session namespaces = main_dbsession.query(Namespace).all() + # Process each namespace separately to avoid long transactions for namespace in namespaces: - process_namespace_scheduled_actions(namespace.id) + try: + process_namespace_scheduled_actions(namespace.id) + except Exception as e: + log.error(f"Error processing namespace {namespace.id}: {e}") log.info("Finished processing scheduled actions") except Exception as e: log.error(f"Error in scheduler: {e}") finally: - release_scheduler_lock(main_dbsession) + release_scheduler_lock() main_dbsession.close() main_engine.dispose() @@ -731,43 +758,9 @@ def process_namespace_scheduled_actions(namespace_id): .all() ) + # Process each action individually with short transactions for action in pending_actions: - try: - media = ( - namespace_dbsession.query(Media) - .filter_by(id=action.media_id) - .first() - ) - if not media: - action.status = "failed" - continue - - if action.action_type == "delete": - namespace_dbsession.delete(media) - log.info(f"Scheduled delete executed for media {media.filename}") - elif action.action_type == "set_public": - media.visibility = "public" - log.info( - f"Scheduled set_public executed for media {media.filename}" - ) - elif action.action_type == "set_private": - media.visibility = "private" - log.info( - f"Scheduled set_private executed for media {media.filename}" - ) - elif action.action_type == "set_unlisted": - media.visibility = "unlisted" - log.info( - f"Scheduled set_unlisted executed for media {media.filename}" - ) - - action.status = "completed" - - except Exception as e: - log.error(f"Error processing scheduled action {action.id}: {e}") - action.status = "failed" - - namespace_dbsession.commit() + process_single_scheduled_action(namespace_dbsession, action) except Exception as e: log.error(f"Error processing namespace {namespace_id}: {e}") @@ -778,6 +771,56 @@ def process_namespace_scheduled_actions(namespace_id): namespace_engine.dispose() +def process_single_scheduled_action(dbsession, action): + """Process a single scheduled action with minimal transaction time.""" + try: + # Start a new transaction for this action + with dbsession.begin(): + # Refresh the action to get latest state + dbsession.refresh(action) + + # Skip if already processed + if action.status != "pending": + return + + media = ( + dbsession.query(Media) + .filter_by(id=action.media_id) + .first() + ) + if not media: + action.status = "failed" + action.completed_at = datetime.datetime.utcnow() + return + + if action.action_type == "delete": + dbsession.delete(media) + log.info(f"Scheduled delete executed for media {media.filename}") + elif action.action_type == "set_public": + media.visibility = "public" + log.info(f"Scheduled set_public executed for media {media.filename}") + elif action.action_type == "set_private": + media.visibility = "private" + log.info(f"Scheduled set_private executed for media {media.filename}") + elif action.action_type == "set_unlisted": + media.visibility = "unlisted" + log.info(f"Scheduled set_unlisted executed for media {media.filename}") + + action.status = "completed" + action.completed_at = datetime.datetime.utcnow() + + except Exception as e: + log.error(f"Error processing scheduled action {action.id}: {e}") + try: + # Try to mark as failed in a separate transaction + with dbsession.begin(): + dbsession.refresh(action) + action.status = "failed" + action.completed_at = datetime.datetime.utcnow() + except Exception as inner_e: + log.error(f"Failed to mark action {action.id} as failed: {inner_e}") + + # Global variables for scheduler thread scheduler_thread = None scheduler_running = False