From 0042f0c36ae751a078b10741e63d322add693ce5 Mon Sep 17 00:00:00 2001 From: dgtlmoon Date: Thu, 22 Jan 2026 10:29:22 +0100 Subject: [PATCH] Memory management improvements for large screenshots, Brotli snapshot improvements (#3798) --- changedetectionio/async_update_worker.py | 32 +++++ .../content_fetchers/playwright.py | 55 +++++--- .../content_fetchers/puppeteer.py | 36 ++++- .../content_fetchers/screenshot_handler.py | 128 ++++++------------ changedetectionio/model/Watch.py | 46 +++++-- requirements.txt | 2 +- 6 files changed, 172 insertions(+), 127 deletions(-) diff --git a/changedetectionio/async_update_worker.py b/changedetectionio/async_update_worker.py index a657868a1..afb59901c 100644 --- a/changedetectionio/async_update_worker.py +++ b/changedetectionio/async_update_worker.py @@ -163,8 +163,10 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec except ProcessorException as e: if e.screenshot: watch.save_screenshot(screenshot=e.screenshot) + e.screenshot = None # Free memory immediately if e.xpath_data: watch.save_xpath_data(data=e.xpath_data) + e.xpath_data = None # Free memory immediately datastore.update_watch(uuid=uuid, update_obj={'last_error': e.message}) process_changedetection_results = False @@ -184,9 +186,11 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec if e.screenshot: watch.save_screenshot(screenshot=e.screenshot, as_error=True) + e.screenshot = None # Free memory immediately if e.xpath_data: watch.save_xpath_data(data=e.xpath_data) + e.xpath_data = None # Free memory immediately process_changedetection_results = False @@ -205,8 +209,10 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec if e.screenshot: watch.save_screenshot(screenshot=e.screenshot, as_error=True) + e.screenshot = None # Free memory immediately if e.xpath_data: watch.save_xpath_data(data=e.xpath_data, as_error=True) + e.xpath_data = None # Free memory immediately if e.page_text: watch.save_error_text(contents=e.page_text) @@ -223,9 +229,11 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec # Filter wasnt found, but we should still update the visual selector so that they can have a chance to set it up again if e.screenshot: watch.save_screenshot(screenshot=e.screenshot) + e.screenshot = None # Free memory immediately if e.xpath_data: watch.save_xpath_data(data=e.xpath_data) + e.xpath_data = None # Free memory immediately # Only when enabled, send the notification if watch.get('filter_failure_notification_send', False): @@ -317,6 +325,7 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec err_text = "Error running JS Actions - Page request - "+e.message if e.screenshot: watch.save_screenshot(screenshot=e.screenshot, as_error=True) + e.screenshot = None # Free memory immediately datastore.update_watch(uuid=uuid, update_obj={'last_error': err_text, 'last_check_status': e.status_code}) process_changedetection_results = False @@ -328,6 +337,7 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec if e.screenshot: watch.save_screenshot(screenshot=e.screenshot, as_error=True) + e.screenshot = None # Free memory immediately datastore.update_watch(uuid=uuid, update_obj={'last_error': err_text, 'last_check_status': e.status_code, @@ -369,9 +379,17 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec if changed_detected or not watch.history_n: if update_handler.screenshot: watch.save_screenshot(screenshot=update_handler.screenshot) + # Free screenshot memory immediately after saving + update_handler.screenshot = None + if hasattr(update_handler, 'fetcher') and hasattr(update_handler.fetcher, 'screenshot'): + update_handler.fetcher.screenshot = None if update_handler.xpath_data: watch.save_xpath_data(data=update_handler.xpath_data) + # Free xpath data memory + update_handler.xpath_data = None + if hasattr(update_handler, 'fetcher') and hasattr(update_handler.fetcher, 'xpath_data'): + update_handler.fetcher.xpath_data = None # Ensure unique timestamp for history if watch.newest_history_key and int(fetch_start_time) == int(watch.newest_history_key): @@ -438,6 +456,20 @@ async def async_update_worker(worker_id, q, notification_q, app, datastore, exec update_handler.fetcher.clear_content() logger.debug(f"Cleared fetcher content for UUID {uuid}") + # Explicitly delete update_handler to free all references + if update_handler: + del update_handler + update_handler = None + + # Force aggressive memory cleanup after clearing + import gc + gc.collect() + try: + import ctypes + ctypes.CDLL('libc.so.6').malloc_trim(0) + except Exception: + pass + except Exception as e: logger.error(f"Worker {worker_id} unexpected error processing {uuid}: {e}") logger.error(f"Worker {worker_id} traceback:", exc_info=True) diff --git a/changedetectionio/content_fetchers/playwright.py b/changedetectionio/content_fetchers/playwright.py index 1bb625f47..e80677b79 100644 --- a/changedetectionio/content_fetchers/playwright.py +++ b/changedetectionio/content_fetchers/playwright.py @@ -13,7 +13,6 @@ from changedetectionio.content_fetchers.exceptions import PageUnloadable, Non200 async def capture_full_page_async(page, screenshot_format='JPEG', watch_uuid=None, lock_viewport_elements=False): import os import time - import multiprocessing start = time.time() watch_info = f"[{watch_uuid}] " if watch_uuid else "" @@ -105,24 +104,29 @@ async def capture_full_page_async(page, screenshot_format='JPEG', watch_uuid=Non stitch_start = time.time() logger.debug(f"{watch_info}Starting stitching of {len(screenshot_chunks)} chunks") - # For small number of chunks (2-3), stitch inline to avoid multiprocessing overhead - # Only use separate process for many chunks (4+) to avoid blocking the event loop - if len(screenshot_chunks) <= 3: - from changedetectionio.content_fetchers.screenshot_handler import stitch_images_inline - screenshot = stitch_images_inline(screenshot_chunks, page_height, SCREENSHOT_MAX_TOTAL_HEIGHT) - else: - # Use separate process for many chunks to avoid blocking - # Always use spawn for thread safety - consistent behavior in tests and production - from changedetectionio.content_fetchers.screenshot_handler import stitch_images_worker - ctx = multiprocessing.get_context('spawn') - parent_conn, child_conn = ctx.Pipe() - p = ctx.Process(target=stitch_images_worker, args=(child_conn, screenshot_chunks, page_height, SCREENSHOT_MAX_TOTAL_HEIGHT)) - p.start() - screenshot = parent_conn.recv_bytes() - p.join() - # Explicit cleanup - del p - del parent_conn, child_conn + # Always use spawn subprocess for ANY stitching (2+ chunks) + # PIL allocates at C level and Python GC never releases it - subprocess exit forces OS to reclaim + # Trade-off: 35MB resource_tracker vs 500MB+ PIL leak in main process + from changedetectionio.content_fetchers.screenshot_handler import stitch_images_worker_raw_bytes + import multiprocessing + import struct + + ctx = multiprocessing.get_context('spawn') + parent_conn, child_conn = ctx.Pipe() + p = ctx.Process(target=stitch_images_worker_raw_bytes, args=(child_conn, page_height, SCREENSHOT_MAX_TOTAL_HEIGHT)) + p.start() + + # Send via raw bytes (no pickle) + parent_conn.send_bytes(struct.pack('I', len(screenshot_chunks))) + for chunk in screenshot_chunks: + parent_conn.send_bytes(chunk) + + screenshot = parent_conn.recv_bytes() + p.join() + + parent_conn.close() + child_conn.close() + del p, parent_conn, child_conn stitch_time = time.time() - stitch_start total_time = time.time() - start @@ -130,9 +134,6 @@ async def capture_full_page_async(page, screenshot_format='JPEG', watch_uuid=Non logger.debug( f"{watch_info}Screenshot complete - Page height: {page_height}px, Capture height: {SCREENSHOT_MAX_TOTAL_HEIGHT}px | " f"Setup: {setup_time:.2f}s, Capture: {capture_time:.2f}s, Stitching: {stitch_time:.2f}s, Total: {total_time:.2f}s") - # Explicit cleanup - del screenshot_chunks - screenshot_chunks = None return screenshot total_time = time.time() - start @@ -403,6 +404,16 @@ class fetcher(Fetcher): # The actual screenshot - this always base64 and needs decoding! horrible! huge CPU usage self.screenshot = await capture_full_page_async(page=self.page, screenshot_format=self.screenshot_format, watch_uuid=watch_uuid, lock_viewport_elements=self.lock_viewport_elements) + # Force aggressive memory cleanup - screenshots are large and base64 decode creates temporary buffers + await self.page.request_gc() + gc.collect() + # Release C-level memory from base64 decode back to OS + try: + import ctypes + ctypes.CDLL('libc.so.6').malloc_trim(0) + except Exception: + pass + except ScreenshotUnavailable: # Re-raise screenshot unavailable exceptions raise diff --git a/changedetectionio/content_fetchers/puppeteer.py b/changedetectionio/content_fetchers/puppeteer.py index a120850ac..e2e9e9410 100644 --- a/changedetectionio/content_fetchers/puppeteer.py +++ b/changedetectionio/content_fetchers/puppeteer.py @@ -23,7 +23,6 @@ from changedetectionio.content_fetchers.exceptions import PageUnloadable, Non200 async def capture_full_page(page, screenshot_format='JPEG', watch_uuid=None, lock_viewport_elements=False): import os import time - import multiprocessing start = time.time() watch_info = f"[{watch_uuid}] " if watch_uuid else "" @@ -122,24 +121,39 @@ async def capture_full_page(page, screenshot_format='JPEG', watch_uuid=None, loc logger.debug(f"{watch_info}All {len(screenshot_chunks)} chunks captured in {capture_time:.2f}s (total chunk time: {total_capture_time:.2f}s)") if len(screenshot_chunks) > 1: - # Always use spawn for thread safety - consistent behavior in tests and production - from changedetectionio.content_fetchers.screenshot_handler import stitch_images_worker stitch_start = time.time() logger.debug(f"{watch_info}Starting stitching of {len(screenshot_chunks)} chunks") + + # Always use spawn subprocess for ANY stitching (2+ chunks) + # PIL allocates at C level and Python GC never releases it - subprocess exit forces OS to reclaim + # Trade-off: 35MB resource_tracker vs 500MB+ PIL leak in main process + from changedetectionio.content_fetchers.screenshot_handler import stitch_images_worker_raw_bytes + import multiprocessing + import struct + ctx = multiprocessing.get_context('spawn') parent_conn, child_conn = ctx.Pipe() - p = ctx.Process(target=stitch_images_worker, args=(child_conn, screenshot_chunks, page_height, SCREENSHOT_MAX_TOTAL_HEIGHT)) + p = ctx.Process(target=stitch_images_worker_raw_bytes, args=(child_conn, page_height, SCREENSHOT_MAX_TOTAL_HEIGHT)) p.start() + + # Send via raw bytes (no pickle) + parent_conn.send_bytes(struct.pack('I', len(screenshot_chunks))) + for chunk in screenshot_chunks: + parent_conn.send_bytes(chunk) + screenshot = parent_conn.recv_bytes() p.join() + + parent_conn.close() + child_conn.close() + del p, parent_conn, child_conn + stitch_time = time.time() - stitch_start total_time = time.time() - start setup_time = total_time - capture_time - stitch_time logger.debug( f"{watch_info}Screenshot complete - Page height: {page_height}px, Capture height: {SCREENSHOT_MAX_TOTAL_HEIGHT}px | " f"Setup: {setup_time:.2f}s, Capture: {capture_time:.2f}s, Stitching: {stitch_time:.2f}s, Total: {total_time:.2f}s") - - screenshot_chunks = None return screenshot total_time = time.time() - start @@ -422,6 +436,16 @@ class fetcher(Fetcher): # Now take screenshot (scrolling may trigger layout changes, but measurements are already captured) logger.debug(f"Screenshot format {self.screenshot_format}") self.screenshot = await capture_full_page(page=self.page, screenshot_format=self.screenshot_format, watch_uuid=watch_uuid, lock_viewport_elements=self.lock_viewport_elements) + + # Force aggressive memory cleanup - pyppeteer base64 decode creates temporary buffers + import gc + gc.collect() + # Release C-level memory from base64 decode back to OS + try: + import ctypes + ctypes.CDLL('libc.so.6').malloc_trim(0) + except Exception: + pass self.xpath_data = await self.page.evaluate(XPATH_ELEMENT_JS, { "visualselector_xpath_selectors": visualselector_xpath_selectors, "max_height": MAX_TOTAL_HEIGHT diff --git a/changedetectionio/content_fetchers/screenshot_handler.py b/changedetectionio/content_fetchers/screenshot_handler.py index f0f765207..fb09f9aeb 100644 --- a/changedetectionio/content_fetchers/screenshot_handler.py +++ b/changedetectionio/content_fetchers/screenshot_handler.py @@ -8,92 +8,42 @@ from loguru import logger from changedetectionio.content_fetchers import SCREENSHOT_MAX_HEIGHT_DEFAULT, SCREENSHOT_DEFAULT_QUALITY -# Cache font to avoid loading on every stitch -_cached_font = None - -def _get_caption_font(): - """Get or create cached font for caption text.""" - global _cached_font - if _cached_font is None: - from PIL import ImageFont - try: - _cached_font = ImageFont.truetype("arial.ttf", 35) - except IOError: - _cached_font = ImageFont.load_default() - return _cached_font - - -def stitch_images_inline(chunks_bytes, original_page_height, capture_height): - """ - Stitch image chunks together inline (no multiprocessing). - Optimized for small number of chunks (2-3) to avoid process creation overhead. - - Args: - chunks_bytes: List of JPEG image bytes - original_page_height: Original page height in pixels - capture_height: Maximum capture height - - Returns: - bytes: Stitched JPEG image - """ - import os - import io - from PIL import Image, ImageDraw - - # Load images from byte chunks - images = [Image.open(io.BytesIO(b)) for b in chunks_bytes] - total_height = sum(im.height for im in images) - max_width = max(im.width for im in images) - - # Create stitched image - stitched = Image.new('RGB', (max_width, total_height)) - y_offset = 0 - for im in images: - stitched.paste(im, (0, y_offset)) - y_offset += im.height - im.close() # Close immediately after pasting - - # Draw caption only if page was trimmed - if original_page_height > capture_height: - draw = ImageDraw.Draw(stitched) - caption_text = f"WARNING: Screenshot was {original_page_height}px but trimmed to {capture_height}px because it was too long" - padding = 10 - font = _get_caption_font() - - bbox = draw.textbbox((0, 0), caption_text, font=font) - text_width = bbox[2] - bbox[0] - text_height = bbox[3] - bbox[1] - - # Draw white background rectangle - draw.rectangle([(0, 0), (max_width, text_height + 2 * padding)], fill=(255, 255, 255)) - - # Draw text centered - text_x = (max_width - text_width) // 2 - draw.text((text_x, padding), caption_text, font=font, fill=(255, 0, 0)) - - # Encode to JPEG - output = io.BytesIO() - stitched.save(output, format="JPEG", quality=int(os.getenv("SCREENSHOT_QUALITY", SCREENSHOT_DEFAULT_QUALITY)), optimize=True) - result = output.getvalue() - - # Cleanup - stitched.close() - - return result - - -def stitch_images_worker(pipe_conn, chunks_bytes, original_page_height, capture_height): +def stitch_images_worker_raw_bytes(pipe_conn, original_page_height, capture_height): """ Stitch image chunks together in a separate process. - Used for large number of chunks (4+) to avoid blocking the main event loop. + + Uses spawn multiprocessing to isolate PIL's C-level memory allocation. + When the subprocess exits, the OS reclaims ALL memory including C-level allocations + that Python's GC cannot release. This prevents the ~50MB per stitch from accumulating + in the main process. + + Trade-off: Adds 35MB resource_tracker subprocess, but prevents 500MB+ memory leak + in main process (much better at scale: 35GB vs 500GB for 1000 instances). + + Args: + pipe_conn: Pipe connection to receive data and send result + original_page_height: Original page height in pixels + capture_height: Maximum capture height """ import os import io + import struct from PIL import Image, ImageDraw, ImageFont try: + # Receive chunk count as 4-byte integer (no pickle!) + count_bytes = pipe_conn.recv_bytes() + chunk_count = struct.unpack('I', count_bytes)[0] + + # Receive each chunk as raw bytes (no pickle!) + chunks_bytes = [] + for _ in range(chunk_count): + chunks_bytes.append(pipe_conn.recv_bytes()) + # Load images from byte chunks images = [Image.open(io.BytesIO(b)) for b in chunks_bytes] + del chunks_bytes + total_height = sum(im.height for im in images) max_width = max(im.width for im in images) @@ -103,15 +53,14 @@ def stitch_images_worker(pipe_conn, chunks_bytes, original_page_height, capture_ for im in images: stitched.paste(im, (0, y_offset)) y_offset += im.height - im.close() # Close immediately after pasting + im.close() + del images # Draw caption only if page was trimmed if original_page_height > capture_height: draw = ImageDraw.Draw(stitched) caption_text = f"WARNING: Screenshot was {original_page_height}px but trimmed to {capture_height}px because it was too long" padding = 10 - - # Try to load font try: font = ImageFont.truetype("arial.ttf", 35) except IOError: @@ -120,23 +69,26 @@ def stitch_images_worker(pipe_conn, chunks_bytes, original_page_height, capture_ bbox = draw.textbbox((0, 0), caption_text, font=font) text_width = bbox[2] - bbox[0] text_height = bbox[3] - bbox[1] - - # Draw white background rectangle draw.rectangle([(0, 0), (max_width, text_height + 2 * padding)], fill=(255, 255, 255)) - - # Draw text centered text_x = (max_width - text_width) // 2 draw.text((text_x, padding), caption_text, font=font, fill=(255, 0, 0)) - # Encode and send image with optimization + # Encode and send output = io.BytesIO() stitched.save(output, format="JPEG", quality=int(os.getenv("SCREENSHOT_QUALITY", SCREENSHOT_DEFAULT_QUALITY)), optimize=True) - pipe_conn.send_bytes(output.getvalue()) + result_bytes = output.getvalue() stitched.close() + del stitched + output.close() + del output + + pipe_conn.send_bytes(result_bytes) + del result_bytes + except Exception as e: - pipe_conn.send(f"error:{e}") + logger.error(f"Error in stitch_images_worker_raw_bytes: {e}") + error_msg = f"error:{e}".encode('utf-8') + pipe_conn.send_bytes(error_msg) finally: pipe_conn.close() - - diff --git a/changedetectionio/model/Watch.py b/changedetectionio/model/Watch.py index f5bc8593b..6eee3e36d 100644 --- a/changedetectionio/model/Watch.py +++ b/changedetectionio/model/Watch.py @@ -20,8 +20,9 @@ mtable = {'seconds': 1, 'minutes': 60, 'hours': 3600, 'days': 86400, 'weeks': 86 def _brotli_save(contents, filepath, mode=None, fallback_uncompressed=False): """ - Save compressed data using native brotli. - Testing shows no memory leak when using gc.collect() after compression. + Save compressed data using native brotli with streaming compression. + Uses chunked compression to minimize peak memory usage and malloc_trim() + to force release of C-level memory back to the OS. Args: contents: data to compress (str or bytes) @@ -37,27 +38,52 @@ def _brotli_save(contents, filepath, mode=None, fallback_uncompressed=False): """ import brotli import gc + import ctypes # Ensure contents are bytes if isinstance(contents, str): contents = contents.encode('utf-8') try: - logger.debug(f"Starting brotli compression of {len(contents)} bytes.") + original_size = len(contents) + logger.debug(f"Starting brotli streaming compression of {original_size} bytes.") - if mode is not None: - compressed_data = brotli.compress(contents, mode=mode) - else: - compressed_data = brotli.compress(contents) + # Create streaming compressor + compressor = brotli.Compressor(quality=6, mode=mode if mode is not None else brotli.MODE_GENERIC) + + # Stream compress in chunks to minimize memory usage + chunk_size = 65536 # 64KB chunks + total_compressed_size = 0 with open(filepath, 'wb') as f: - f.write(compressed_data) + # Process data in chunks + offset = 0 + while offset < len(contents): + chunk = contents[offset:offset + chunk_size] + compressed_chunk = compressor.process(chunk) + if compressed_chunk: + f.write(compressed_chunk) + total_compressed_size += len(compressed_chunk) + offset += chunk_size - logger.debug(f"Finished brotli compression - From {len(contents)} to {len(compressed_data)} bytes.") + # Finalize compression - critical for proper cleanup + final_chunk = compressor.finish() + if final_chunk: + f.write(final_chunk) + total_compressed_size += len(final_chunk) - # Force garbage collection to prevent memory buildup + logger.debug(f"Finished brotli compression - From {original_size} to {total_compressed_size} bytes.") + + # Cleanup: Delete compressor, force Python GC, then force C-level memory release + del compressor gc.collect() + # Force release of C-level memory back to OS (since brotli is a C library) + try: + ctypes.CDLL('libc.so.6').malloc_trim(0) + except Exception: + pass # malloc_trim not available on all systems (e.g., macOS) + return filepath except Exception as e: diff --git a/requirements.txt b/requirements.txt index 420957ece..6ba18e0a3 100644 --- a/requirements.txt +++ b/requirements.txt @@ -91,7 +91,7 @@ jq~=1.3; python_version >= "3.8" and sys_platform == "linux" # playwright is installed at Dockerfile build time because it's not available on all platforms -pyppeteer-ng==2.0.0rc11 +pyppeteer-ng==2.0.0rc12 pyppeteerstealth>=0.0.4 # Include pytest, so if theres a support issue we can ask them to run these tests on their setup