From 41321889bb45463880cedd1c83a211c9f4a2972e Mon Sep 17 00:00:00 2001 From: dgtlmoon Date: Tue, 13 Jan 2026 13:52:21 +0100 Subject: [PATCH] Async worker updates and increase testing --- changedetectionio/async_update_worker.py | 19 +- changedetectionio/queue_handlers.py | 6 +- changedetectionio/tests/conftest.py | 45 +++ .../tests/test_history_consistency.py | 28 +- changedetectionio/tests/util.py | 2 + changedetectionio/worker_handler.py | 330 ++++++++---------- 6 files changed, 232 insertions(+), 198 deletions(-) diff --git a/changedetectionio/async_update_worker.py b/changedetectionio/async_update_worker.py index 4079cf9ea..388d4da63 100644 --- a/changedetectionio/async_update_worker.py +++ b/changedetectionio/async_update_worker.py @@ -32,21 +32,32 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore): task = asyncio.current_task() if task: task.set_name(f"async-worker-{worker_id}") - + logger.info(f"Starting async worker {worker_id}") - + while not app.config.exit.is_set(): update_handler = None watch = None try: - # Use native janus async interface - no threads needed! - queued_item_data = await asyncio.wait_for(q.async_get(), timeout=1.0) + # Use sync interface via run_in_executor since each worker has its own event loop + loop = asyncio.get_event_loop() + queued_item_data = await asyncio.wait_for( + loop.run_in_executor(None, q.get, True, 1.0), # block=True, timeout=1.0 + timeout=1.5 + ) except asyncio.TimeoutError: # No jobs available, continue loop continue except Exception as e: + # Handle expected Empty exception from queue timeout + import queue + if isinstance(e, queue.Empty): + # Queue is empty, normal behavior - just continue + continue + + # Unexpected exception - log as critical logger.critical(f"CRITICAL: Worker {worker_id} failed to get queue item: {type(e).__name__}: {e}") # Log queue health for debugging diff --git a/changedetectionio/queue_handlers.py b/changedetectionio/queue_handlers.py index b13f6b21c..8ddaee159 100644 --- a/changedetectionio/queue_handlers.py +++ b/changedetectionio/queue_handlers.py @@ -86,6 +86,7 @@ class RecheckPriorityQueue: def get(self, block: bool = True, timeout: Optional[float] = None): """Thread-safe sync get with priority ordering""" + import queue try: # Wait for notification self.sync_q.get(block=block, timeout=timeout) @@ -103,8 +104,11 @@ class RecheckPriorityQueue: logger.debug(f"Successfully retrieved item: {self._get_item_uuid(item)}") return item + except queue.Empty: + # Queue is empty with timeout - expected behavior, re-raise without logging + raise except Exception as e: - logger.critical(f"CRITICAL: Failed to get item from queue: {str(e)}") + # Re-raise without logging - caller (worker) will handle and log appropriately raise # ASYNC INTERFACE (for workers) diff --git a/changedetectionio/tests/conftest.py b/changedetectionio/tests/conftest.py index 24984a270..2583becbd 100644 --- a/changedetectionio/tests/conftest.py +++ b/changedetectionio/tests/conftest.py @@ -270,3 +270,48 @@ def app(request, datastore_path): request.addfinalizer(teardown) yield app + + +@pytest.fixture(scope='session') +def live_server(app): + """Create a live server with explicit threading enabled. + + While pytest-flask's default uses threaded=True, using Werkzeug's + make_server directly with threading appears to be faster for + concurrent test requests, likely due to avoiding multiprocessing overhead. + """ + from werkzeug.serving import make_server + import socket + + # Find a free port + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: + s.bind(('', 0)) + s.listen(1) + port = s.getsockname()[1] + + # Save original SERVER_NAME and set it to include the port + # This is critical for url_for(..., _external=True) to work correctly + original_server_name = app.config.get("SERVER_NAME") or "localhost" + app.config["SERVER_NAME"] = f"127.0.0.1:{port}" + + # Create threaded server + server = make_server('127.0.0.1', port, app, threaded=True) + thread = Thread(target=server.serve_forever, daemon=True) + thread.start() + + # Add attributes for compatibility with pytest-flask's LiveServer + server.url = lambda url='': f'http://127.0.0.1:{port}{url}' + server.app = app + server.host = '127.0.0.1' + server.port = port + + # Add dummy start() method since server is already started + # (called by new_live_server_setup in util.py:338) + server.start = lambda: None + + yield server + + server.shutdown() + + # Restore original SERVER_NAME + app.config["SERVER_NAME"] = original_server_name diff --git a/changedetectionio/tests/test_history_consistency.py b/changedetectionio/tests/test_history_consistency.py index 75ec2b1c7..9191b7632 100644 --- a/changedetectionio/tests/test_history_consistency.py +++ b/changedetectionio/tests/test_history_consistency.py @@ -4,28 +4,27 @@ import time import os import json from flask import url_for - +from loguru import logger from build.lib.changedetectionio.strtobool import strtobool from .util import wait_for_all_checks, delete_all_watches -from urllib.parse import urlparse, parse_qs +import brotli def test_consistent_history(client, live_server, measure_memory_usage, datastore_path): - # live_server_setup(live_server) # Setup on conftest per function - workers = int(os.getenv("FETCH_WORKERS", 10)) - r = range(1, 10 + workers) - uuids = set() - import brotli - for one in r: + uuids = set() + workers = range(1, int(os.getenv("FETCH_WORKERS", 10))) + now = time.time() + + for one in workers: if strtobool(os.getenv("TEST_WITH_BROTLI")): # A very long string that WILL trigger Brotli compression of the snapshot # BROTLI_COMPRESS_SIZE_THRESHOLD should be set to say 200 from ..model.Watch import BROTLI_COMPRESS_SIZE_THRESHOLD - content = str(one) + " " + str(one) * (BROTLI_COMPRESS_SIZE_THRESHOLD + 10) + content = str(one) + "x" + str(one) * (BROTLI_COMPRESS_SIZE_THRESHOLD + 10) else: # Just enough to test datastore - content = str(one) + content = str(one)+'x' test_url = url_for('test_endpoint', content_type="text/html", content=content, _external=True) uuids.add(client.application.config.get('DATASTORE').add_watch(url=test_url, extras={'title': str(one)})) @@ -33,6 +32,9 @@ def test_consistent_history(client, live_server, measure_memory_usage, datastore client.get(url_for("ui.form_watch_checknow"), follow_redirects=True) wait_for_all_checks(client) + logger.debug(f"All fetched in {time.time() - now:.2f}s") + + assert time.time() - now < 10, "Time to fetch all the items should be less than 10 sec, or it could mean something is blocking async workers or other." # Essentially just triggers the DB write/update res = client.post( @@ -44,7 +46,7 @@ def test_consistent_history(client, live_server, measure_memory_usage, datastore ) assert b"Settings updated." in res.data - + # Wait for the sync DB save to happen time.sleep(2) json_db_file = os.path.join(live_server.app.config['DATASTORE'].datastore_path, 'url-watches.json') @@ -54,7 +56,7 @@ def test_consistent_history(client, live_server, measure_memory_usage, datastore json_obj = json.load(f) # assert the right amount of watches was found in the JSON - assert len(json_obj['watching']) == len(r), "Correct number of watches was found in the JSON" + assert len(json_obj['watching']) == len(workers), "Correct number of watches was found in the JSON" i = 0 # each one should have a history.txt containing just one line @@ -91,7 +93,7 @@ def test_consistent_history(client, live_server, measure_memory_usage, datastore watch_title = json_obj['watching'][w]['title'] assert json_obj['watching'][w]['title'], "Watch should have a title set" - assert contents.startswith(watch_title + " "), f"Snapshot file {fname} should start with '{watch_title} '" + assert contents.startswith(watch_title + "x"), f"Snapshot contents in file {fname} should start with '{watch_title}x', got '{contents}'" assert len(files_in_watch_dir) == 3, "Should be just three files in the dir, html.br snapshot, history.txt and the extracted text snapshot" diff --git a/changedetectionio/tests/util.py b/changedetectionio/tests/util.py index 4fe075002..7b7ca936b 100644 --- a/changedetectionio/tests/util.py +++ b/changedetectionio/tests/util.py @@ -189,6 +189,8 @@ def new_live_server_setup(live_server): @live_server.app.route('/test-endpoint') def test_endpoint(): + from loguru import logger + logger.debug(f"/test-endpoint hit {request}") ctype = request.args.get('content_type') status_code = request.args.get('status_code') content = request.args.get('content') or None diff --git a/changedetectionio/worker_handler.py b/changedetectionio/worker_handler.py index cd81007aa..6c98ead25 100644 --- a/changedetectionio/worker_handler.py +++ b/changedetectionio/worker_handler.py @@ -2,7 +2,7 @@ Worker management module for changedetection.io Handles asynchronous workers for dynamic worker scaling. -Sync worker support has been removed in favor of async-only architecture. +Each worker runs in its own thread with its own event loop for isolation. """ import asyncio @@ -11,10 +11,8 @@ import threading import time from loguru import logger -# Global worker state -running_async_tasks = [] -async_loop = None -async_loop_thread = None +# Global worker state - each worker has its own thread and event loop +worker_threads = [] # List of WorkerThread objects # Track currently processing UUIDs for async workers - maps {uuid: worker_id} currently_processing_uuids = {} @@ -23,71 +21,89 @@ currently_processing_uuids = {} USE_ASYNC_WORKERS = True -def start_async_event_loop(): - """Start a dedicated event loop for async workers in a separate thread""" - global async_loop - logger.info("Starting async event loop for workers") - - try: - # Create a new event loop for this thread - async_loop = asyncio.new_event_loop() - # Set it as the event loop for this thread - asyncio.set_event_loop(async_loop) - - logger.debug(f"Event loop created and set: {async_loop}") - - # Run the event loop forever - async_loop.run_forever() - except Exception as e: - logger.error(f"Async event loop error: {e}") - finally: - # Clean up - if async_loop and not async_loop.is_closed(): - async_loop.close() - async_loop = None - logger.info("Async event loop stopped") +class WorkerThread: + """Container for a worker thread with its own event loop""" + def __init__(self, worker_id, update_q, notification_q, app, datastore): + self.worker_id = worker_id + self.update_q = update_q + self.notification_q = notification_q + self.app = app + self.datastore = datastore + self.thread = None + self.loop = None + self.running = False + + def run(self): + """Run the worker in its own event loop""" + try: + # Create a new event loop for this thread + self.loop = asyncio.new_event_loop() + asyncio.set_event_loop(self.loop) + self.running = True + + # Run the worker coroutine + self.loop.run_until_complete( + start_single_async_worker( + self.worker_id, + self.update_q, + self.notification_q, + self.app, + self.datastore + ) + ) + except asyncio.CancelledError: + # Normal shutdown - worker was cancelled + import os + in_pytest = "pytest" in os.sys.modules or "PYTEST_CURRENT_TEST" in os.environ + if not in_pytest: + logger.info(f"Worker {self.worker_id} shutting down gracefully") + except RuntimeError as e: + # Ignore expected shutdown errors + if "Event loop stopped" not in str(e) and "Event loop is closed" not in str(e): + logger.error(f"Worker {self.worker_id} runtime error: {e}") + except Exception as e: + logger.error(f"Worker {self.worker_id} thread error: {e}") + finally: + # Clean up + if self.loop and not self.loop.is_closed(): + self.loop.close() + self.running = False + self.loop = None + + def start(self): + """Start the worker thread""" + self.thread = threading.Thread(target=self.run, daemon=True, name=f"Worker-{self.worker_id}") + self.thread.start() + + def stop(self): + """Stop the worker thread""" + if self.loop and self.running: + try: + # Signal the loop to stop + self.loop.call_soon_threadsafe(self.loop.stop) + except RuntimeError: + pass + + if self.thread and self.thread.is_alive(): + self.thread.join(timeout=2.0) def start_async_workers(n_workers, update_q, notification_q, app, datastore): - """Start the async worker management system""" - global async_loop_thread, async_loop, running_async_tasks, currently_processing_uuids - + """Start async workers, each with its own thread and event loop for isolation""" + global worker_threads, currently_processing_uuids + # Clear any stale UUID tracking state currently_processing_uuids.clear() - - # Start the event loop in a separate thread - async_loop_thread = threading.Thread(target=start_async_event_loop, daemon=True) - async_loop_thread.start() - - # Wait for the loop to be available (with timeout for safety) - max_wait_time = 5.0 - wait_start = time.time() - while async_loop is None and (time.time() - wait_start) < max_wait_time: - time.sleep(0.1) - - if async_loop is None: - logger.error("Failed to start async event loop within timeout") - return - - # Additional brief wait to ensure loop is running - time.sleep(0.2) - - # Start async workers - logger.info(f"Starting {n_workers} async workers") + + # Start each worker in its own thread with its own event loop + logger.info(f"Starting {n_workers} async workers (isolated threads)") for i in range(n_workers): try: - # Use a factory function to create named worker coroutines - def create_named_worker(worker_id): - async def named_worker(): - task = asyncio.current_task() - if task: - task.set_name(f"async-worker-{worker_id}") - return await start_single_async_worker(worker_id, update_q, notification_q, app, datastore) - return named_worker() - - task_future = asyncio.run_coroutine_threadsafe(create_named_worker(i), async_loop) - running_async_tasks.append(task_future) - except RuntimeError as e: + worker = WorkerThread(i, update_q, notification_q, app, datastore) + worker.start() + worker_threads.append(worker) + # No sleep needed - threads start independently and asynchronously + except Exception as e: logger.error(f"Failed to start async worker {i}: {e}") continue @@ -102,8 +118,6 @@ async def start_single_async_worker(worker_id, update_q, notification_q, app, da while not app.config.exit.is_set(): try: - if not in_pytest: - logger.info(f"Starting async worker {worker_id}") await async_update_worker(worker_id, update_q, notification_q, app, datastore) # If we reach here, worker exited cleanly if not in_pytest: @@ -131,39 +145,38 @@ def start_workers(n_workers, update_q, notification_q, app, datastore): def add_worker(update_q, notification_q, app, datastore): """Add a new async worker (for dynamic scaling)""" - global running_async_tasks - - if not async_loop: - logger.error("Async loop not running, cannot add worker") - return False - - worker_id = len(running_async_tasks) + global worker_threads + + worker_id = len(worker_threads) logger.info(f"Adding async worker {worker_id}") - - task_future = asyncio.run_coroutine_threadsafe( - start_single_async_worker(worker_id, update_q, notification_q, app, datastore), async_loop - ) - running_async_tasks.append(task_future) - return True + + try: + worker = WorkerThread(worker_id, update_q, notification_q, app, datastore) + worker.start() + worker_threads.append(worker) + return True + except Exception as e: + logger.error(f"Failed to add worker {worker_id}: {e}") + return False def remove_worker(): """Remove an async worker (for dynamic scaling)""" - global running_async_tasks - - if not running_async_tasks: + global worker_threads + + if not worker_threads: return False - - # Cancel the last worker - task_future = running_async_tasks.pop() - task_future.cancel() - logger.info(f"Removed async worker, {len(running_async_tasks)} workers remaining") + + # Stop the last worker + worker = worker_threads.pop() + worker.stop() + logger.info(f"Removed async worker, {len(worker_threads)} workers remaining") return True def get_worker_count(): """Get current number of async workers""" - return len(running_async_tasks) + return len(worker_threads) def get_running_uuids(): @@ -249,38 +262,21 @@ def queue_item_async_safe(update_q, item, silent=False): def shutdown_workers(): """Shutdown all async workers fast and aggressively""" - global async_loop, async_loop_thread, running_async_tasks - + global worker_threads + # Check if we're in pytest environment - if so, be more gentle with logging import os in_pytest = "pytest" in os.sys.modules or "PYTEST_CURRENT_TEST" in os.environ - + if not in_pytest: logger.info("Fast shutdown of async workers initiated...") - - # Cancel all async tasks immediately - for task_future in running_async_tasks: - if not task_future.done(): - task_future.cancel() - - # Stop the async event loop immediately - if async_loop and not async_loop.is_closed(): - try: - async_loop.call_soon_threadsafe(async_loop.stop) - except RuntimeError: - # Loop might already be stopped - pass - - running_async_tasks.clear() - async_loop = None - - # Give async thread minimal time to finish, then continue - if async_loop_thread and async_loop_thread.is_alive(): - async_loop_thread.join(timeout=1.0) # Only 1 second timeout - if async_loop_thread.is_alive() and not in_pytest: - logger.info("Async thread still running after timeout - continuing with shutdown") - async_loop_thread = None - + + # Stop all worker threads + for worker in worker_threads: + worker.stop() + + worker_threads.clear() + if not in_pytest: logger.info("Async workers fast shutdown complete") @@ -290,69 +286,57 @@ def shutdown_workers(): def adjust_async_worker_count(new_count, update_q=None, notification_q=None, app=None, datastore=None): """ Dynamically adjust the number of async workers. - + Args: new_count: Target number of workers update_q, notification_q, app, datastore: Required for adding new workers - + Returns: dict: Status of the adjustment operation """ - global running_async_tasks - + global worker_threads + current_count = get_worker_count() - + if new_count == current_count: return { 'status': 'no_change', 'message': f'Worker count already at {current_count}', 'current_count': current_count } - + if new_count > current_count: # Add workers workers_to_add = new_count - current_count logger.info(f"Adding {workers_to_add} async workers (from {current_count} to {new_count})") - + if not all([update_q, notification_q, app, datastore]): return { 'status': 'error', 'message': 'Missing required parameters to add workers', 'current_count': current_count } - + for i in range(workers_to_add): - worker_id = len(running_async_tasks) - task_future = asyncio.run_coroutine_threadsafe( - start_single_async_worker(worker_id, update_q, notification_q, app, datastore), - async_loop - ) - running_async_tasks.append(task_future) - + add_worker(update_q, notification_q, app, datastore) + return { 'status': 'success', 'message': f'Added {workers_to_add} workers', 'previous_count': current_count, - 'current_count': new_count + 'current_count': len(worker_threads) } - + else: # Remove workers workers_to_remove = current_count - new_count logger.info(f"Removing {workers_to_remove} async workers (from {current_count} to {new_count})") - + removed_count = 0 for _ in range(workers_to_remove): - if running_async_tasks: - task_future = running_async_tasks.pop() - task_future.cancel() - # Wait for the task to actually stop - try: - task_future.result(timeout=5) # 5 second timeout - except Exception: - pass # Task was cancelled, which is expected + if remove_worker(): removed_count += 1 - + return { 'status': 'success', 'message': f'Removed {removed_count} workers', @@ -367,72 +351,58 @@ def get_worker_status(): 'worker_type': 'async', 'worker_count': get_worker_count(), 'running_uuids': get_running_uuids(), - 'async_loop_running': async_loop is not None, + 'active_threads': sum(1 for w in worker_threads if w.thread and w.thread.is_alive()), } def check_worker_health(expected_count, update_q=None, notification_q=None, app=None, datastore=None): """ Check if the expected number of async workers are running and restart any missing ones. - + Args: expected_count: Expected number of workers update_q, notification_q, app, datastore: Required for restarting workers - + Returns: dict: Health check results """ - global running_async_tasks - + global worker_threads + current_count = get_worker_count() - - if current_count == expected_count: + + # Check which workers are actually alive + alive_count = sum(1 for w in worker_threads if w.thread and w.thread.is_alive()) + + if alive_count == expected_count: return { 'status': 'healthy', 'expected_count': expected_count, - 'actual_count': current_count, + 'actual_count': alive_count, 'message': f'All {expected_count} async workers running' } - - # Check for crashed async workers + + # Find dead workers dead_workers = [] - alive_count = 0 - - for i, task_future in enumerate(running_async_tasks[:]): - if task_future.done(): - try: - result = task_future.result() - dead_workers.append(i) - logger.warning(f"Async worker {i} completed unexpectedly") - except Exception as e: - dead_workers.append(i) - logger.error(f"Async worker {i} crashed: {e}") - else: - alive_count += 1 - + for i, worker in enumerate(worker_threads[:]): + if not worker.thread or not worker.thread.is_alive(): + dead_workers.append(i) + logger.warning(f"Async worker {worker.worker_id} thread is dead") + # Remove dead workers from tracking for i in reversed(dead_workers): - if i < len(running_async_tasks): - running_async_tasks.pop(i) - + if i < len(worker_threads): + worker_threads.pop(i) + missing_workers = expected_count - alive_count restarted_count = 0 - + if missing_workers > 0 and all([update_q, notification_q, app, datastore]): logger.info(f"Restarting {missing_workers} crashed async workers") - + for i in range(missing_workers): - worker_id = alive_count + i - try: - task_future = asyncio.run_coroutine_threadsafe( - start_single_async_worker(worker_id, update_q, notification_q, app, datastore), - async_loop - ) - running_async_tasks.append(task_future) + if add_worker(update_q, notification_q, app, datastore): restarted_count += 1 - except Exception as e: - logger.error(f"Failed to restart worker {worker_id}: {e}") - + return { 'status': 'repaired' if restarted_count > 0 else 'degraded', 'expected_count': expected_count,