From 895150ead80edb49d58066c001534fb9c068afa8 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Mon, 22 Sep 2025 15:53:01 +0300 Subject: [PATCH 01/30] feat: docker entry point and celery config - Add Docker entrypoint script for Xvfb and x11vnc setup - Add celery configuration - not configured yet. --- celery_config.py | 21 +++++++++++++++++++++ docker-entrypoint.sh | 14 ++++++++++++++ 2 files changed, 35 insertions(+) create mode 100644 celery_config.py create mode 100644 docker-entrypoint.sh diff --git a/celery_config.py b/celery_config.py new file mode 100644 index 0000000..1fa05f0 --- /dev/null +++ b/celery_config.py @@ -0,0 +1,21 @@ +# celeryconfig.py - Celery configuration file + +# Broker settings (LavinMQ) + +# Replace BROKER_URL with your LavinMQ server's URL +# Example : 'amqp://guest:guest@localhost:5672//' +# Default password and username for lavinmq user: guest, pass: guest + +BROKER_URL = 'lavinmq://:@:/' + +# Result backend (Optional, if you want to store task results) +# result_backend = 'rpc:// + +# Recommended settings for local LavinMQ +broker_pool_limit = 1 +broker_heartbeat = None +broker_connection_timeout = 30 +result_backend = None +event_queue_expires = 60 +worker_prefetch_multiplier = 1 +worker_concurrency = 4 # Adjust based on your system's capabilities \ No newline at end of file diff --git a/docker-entrypoint.sh b/docker-entrypoint.sh new file mode 100644 index 0000000..4d1cf56 --- /dev/null +++ b/docker-entrypoint.sh @@ -0,0 +1,14 @@ +#!/bin/sh + +# Ensure the X11 socket directory exists and has correct permissions +# This is done as root before switching to appuser +mkdir -p /tmp/.X11-unix +chmod 1777 /tmp/.X11-unix +chown appuser:appuser /tmp/.X11-unix + +# Start Xvfb and x11vnc in the background +Xvfb :99 -screen 0 1280x720x24 -ac & +x11vnc -display :99 -nopw -forever -shared -rfbport 5900 -quiet & + +# Execute the main command passed to the entrypoint +exec "$@" From 17a70b496895836dc839445f7d2bc46c581b3ed8 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Mon, 22 Sep 2025 15:56:20 +0300 Subject: [PATCH 02/30] feat(db): enhance Firebase and Redis configurations for production and emulator support --- firestore.py | 30 ++++++++++++++++++++++++++---- 1 file changed, 26 insertions(+), 4 deletions(-) diff --git a/firestore.py b/firestore.py index 49df3d7..bfbd737 100644 --- a/firestore.py +++ b/firestore.py @@ -8,14 +8,24 @@ from firebase_admin import credentials import json - +# Determine which .env file to load +production = os.getenv("PYTHON_ENV", "development").lower() == "production" +env_file = ".env" if production else "dev.env" # Load environment variables -load_dotenv() +load_dotenv(env_file) # config = load_config() - fire_config = os.getenv("FIRESTORE_CONFIG") firebase_cred = os.getenv("FIREBASE_CREDENTIALS") +# Emulator settings +print(f"\033[93mINFO-DB:\033[0m Running in {'production' if production else 'development'} mode.") + +if not production: + firestore_emulator_host = os.getenv("FIRESTORE_EMULATOR_HOST") + firebase_auth_emulator_host = os.getenv("FIREBASE_AUTH_EMULATOR_HOST") +else: + firestore_emulator_host = None + firebase_auth_emulator_host = None if not fire_config or not firebase_cred: @@ -37,6 +47,18 @@ def __init__(self): self.company_credentials = credentials.Certificate(self.firebase_credentials) def init_firebase(self): + # Set environment variables for emulators if present + if not production and (firestore_emulator_host or firebase_auth_emulator_host): + if firestore_emulator_host: + os.environ["FIRESTORE_EMULATOR_HOST"] = firestore_emulator_host + print(f"\033[93mINFO-DB:\033[0m Using Firestore emulator at {firestore_emulator_host}") + if firebase_auth_emulator_host: + os.environ["FIREBASE_AUTH_EMULATOR_HOST"] = firebase_auth_emulator_host + print(f"\033[93mINFO-DB:\033[0m Using Firebase Auth emulator at {firebase_auth_emulator_host}") + + print("\033[93mWARNING-DB:\033[0m Running in emulator mode. Data will not be persisted.") + self.company_credentials = credentials.ApplicationDefault() + # # Initialize Firebase if not already initialized if not firebase_admin._apps: # Initialize Firebase app with the company credentials @@ -45,7 +67,7 @@ def init_firebase(self): # Set up Firestore with a custom database ID dbName = firestoreConfig["dbName"] # Replace with your Firestore database ID - self.db = firebase_admin.firestore.client(database_id=dbName) + self.db = firebase_admin.firestore.client(self.app, database_id=dbName) # Assign the authentication object self.auth = firebase_admin.auth From 080403e74efe73d205d8f0449196821fa621c708 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Mon, 22 Sep 2025 15:59:15 +0300 Subject: [PATCH 03/30] chore(config): update fly.toml and requirements for improved process management and dependency versioning - Update Fly configuration for production environment and process management -Update requirements to specify compatible versions - update CORS origins and enhance environment settings for production --- config.yml | 7 ++++++- fly.toml | 8 ++++++-- requirements.txt | 3 ++- 3 files changed, 14 insertions(+), 4 deletions(-) diff --git a/config.yml b/config.yml index b0921f9..77990e6 100644 --- a/config.yml +++ b/config.yml @@ -8,8 +8,13 @@ app: workers: 1 uvloop: auto timeout_keep_alive: 300 - cors_origins: + cors_origins_dev: - "http://localhost:4200" + - "http://localhost:8081" + - "http://localhost:5000" + - "https://deepscrape.web.app" + - "https://deepscrape.dev" + cors_origins: - "https://deepscrape.web.app" - "https://deepscrape.dev" gzip: diff --git a/fly.toml b/fly.toml index fe68f40..72e8b07 100644 --- a/fly.toml +++ b/fly.toml @@ -9,6 +9,7 @@ primary_region = 'fra' [build] [env] + PYTHON_ENV = 'production' UPSTASH_REDIS_PORT = '30766' UPSTASH_REDIS_REST_URL = 'https://gusc1-saved-terrapin-30766.upstash.io' AWS_ENDPOINT_URL_S3 = 'https://fly.storage.tigris.dev' @@ -55,8 +56,10 @@ primary_region = 'fra' handlers = ["tls", "http"] # Ensure "tls" is included for HTTPS [processes] - app = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & x11vnc -display :99 -nopw -forever -shared -rfbport 5900 -quiet & /opt/noVNC/utils/websockify/run 6080 localhost:5900 --web /opt/noVNC & uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets'" - worker = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & x11vnc -display :99 -nopw -forever -shared -rfbport 5900 -quiet & celery -A celery_app.celery_app worker --loglevel=info'" + app = "sh -c '/opt/noVNC/utils/websockify/run --web /opt/noVNC 0.0.0.0:6080 0.0.0.0:5900 & uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets'" + worker = "celery -A celery_app.celery_app worker --loglevel=info" + + # app = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets'" # worker = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && celery -A celery_app.celery_app worker --loglevel=info --concurrency=2'" # app = "uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets" @@ -66,6 +69,7 @@ primary_region = 'fra' memory = '2gb' cpu_kind = 'shared' cpus = 4 + swap_size_mb = 512 # Add swap for compilation [[metrics]] port = 8000 diff --git a/requirements.txt b/requirements.txt index f89f405..a9aff59 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,6 +1,7 @@ fastapi gunicorn dotenv +websockets uvicorn[standard] crawl4ai uvloop @@ -8,7 +9,7 @@ beautifulsoup4~=4.12 tf-playwright-stealth>=1.1.0 firebase-admin google-cloud-firestore -upstash-redis +upstash-redis~=1.4.0 upstash_ratelimit apscheduler pydantic>=2.10 From d8dcec6ae3a1c251383ccf59695ddd2e1adb25ae Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Mon, 22 Sep 2025 16:00:53 +0300 Subject: [PATCH 04/30] feat(redis): enhance Redis connection handling and add stream message functions --- redisCache.py | 104 ++++++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 93 insertions(+), 11 deletions(-) diff --git a/redisCache.py b/redisCache.py index cced7bd..11c9a13 100644 --- a/redisCache.py +++ b/redisCache.py @@ -1,17 +1,21 @@ # for async use -import asyncio import os -from typing import List +from typing import List, Optional from dotenv import load_dotenv # from redis.asyncio import Redis from upstash_ratelimit.asyncio import Ratelimit, FixedWindow, TokenBucket from upstash_redis.asyncio import Redis from redis.asyncio import Redis as PureRedis +import asyncio REDIS_CHANNEL = "stream_channel" # Default channel for streaming data # ────────────────── configuration ────────────────── -load_dotenv() +production = os.getenv("PYTHON_ENV", "development").lower() == "production" +env_file = ".env" if production else "dev.env" + +# Load environment variables +load_dotenv(env_file, verbose=True) redis_url = os.environ.get("UPSTASH_REDIS_REST_URL") redis_token = os.environ.get("UPSTASH_REDIS_REST_TOKEN") @@ -41,16 +45,28 @@ username=REDIS_USERNAME, password=REDIS_PASSWORD, ssl=True, # Use SSL if your Redis server supports it - decode_responses=True # Optional: decode responses to UTF-8 strings + decode_responses=True, # Optional: decode responses to UTF-8 strings ) async def test_connection(redis: Redis | PureRedis): - try: - await redis.ping() - - print("\033[94mINFO-DB:\033[0m \033[92mRedis connected successfully!\033[0m") - except Exception as e: - print(f"\033[91mERROR-DB:\033[0m Redis connection failed: {e}") + retries = 0 + max_retries = 100 + retry_delay = 12 # seconds + + while retries < max_retries: + try: + await redis.ping() + print("\033[94mINFO-DB:\033[0m \033[92mRedis connected successfully!\033[0m") + return + except Exception as e: + retries += 1 + print(f"\033[91mERROR-DB:\033[0m Redis connection failed: {e}") + print(f"\033[93mWARNING-DB:\033[0m Trying again in 12.00 seconds... (attempt {retries}/{max_retries})") + if retries < max_retries: + await asyncio.sleep(retry_delay) + else: + print("\033[91mERROR-DB:\033[0m Max retries reached. Could not connect to Redis.") + return # ─────────────────── rate limiters ────────────────── @@ -106,4 +122,70 @@ async def redis_subscribe(redis: Redis, channel: str): except Exception as e: print(f"\033[91mERROR-DB:\033[0m Failed to subscribe to channel '{channel}': {e}") return None - \ No newline at end of file + + +async def redis_xadd(pipe, channel: str, message: dict, maxlen: Optional[int] = None, approximate: bool = False): + """Add a message to a Redis stream with optional maxlen.""" + if not pipe: + print("\033[91mERROR-DB:\033[0m Redis pipe client is not initialized.") + return None + if not channel: + print("\033[91mERROR-DB:\033[0m Channel is empty.") + return None + if not message: + print("\033[91mERROR-DB:\033[0m Message is empty.") + return None + + try: + pieces = [channel] + if maxlen is not None: + if maxlen < 0: + raise ValueError("maxlen must be a non-negative integer") + pieces.append("MAXLEN") + if approximate: + pieces.append("~") + pieces.append(str(maxlen)) + for key, value in message.items(): + pieces.extend([key, value]) + + message_id = pipe.execute(["XADD", *pieces]) + return message_id + except Exception as e: + print(f"\033[91mERROR-DB:\033[0m Failed to add message to stream '{channel}': {e}") + return None + + +async def redis_xread(redis: Redis, streams: dict, count: Optional[int] = None, block: Optional[int] = None): + """ + Read messages from a Redis stream using XREAD. + :param redis: Redis client (PureRedis) + :param streams: Dictionary of {channel: last_id} + :param count: Maximum number of entries to return + :param block: Number of milliseconds to block if no messages are available + :return: List of messages or None + """ + if not redis: + print("\033[91mERROR-DB:\033[0m Redis redis client is not initialized.") + return None + if not streams: + print("\033[91mERROR-DB:\033[0m Streams dictionary is empty.") + return None + + try: + # Build the XREAD command as a list + command = ["XREAD"] + if count is not None: + command.extend(["COUNT", str(count)]) + if block is not None: + command.extend(["BLOCK", str(block)]) + command.append("STREAMS") + for channel, last_id in streams.items(): + command.append(channel) + for channel, last_id in streams.items(): + command.append(last_id) + result = await redis.execute(command) + return result + except Exception as e: + print(f"\033[91mERROR-DB:\033[0m Failed to read from stream(s) '{list(streams.keys())}': {e}") + return None +# ─────────────────── redis pub/sub ────────────────── From 135c2e15b18947062657e44b6c8c7fe23b604429 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Mon, 22 Sep 2025 16:05:41 +0300 Subject: [PATCH 05/30] feat(server): This commit introduces significant updates to the application's Docker build process, core dependencies, and operational stability - enhance task and job handling with improved event loop management and error handling - Update unsupported Accept header response to use status constant - Improve server.py for better WebSocket client management and logging - Refactor utils.py for better WebSocket client cleanup and error handling - Modify job.py to improve task cancellation and status retrieval logic - Implemented platform-specific asynchronous event loop policies, utilizing `uvloop` for non-Windows systems and `WindowsProactorEventLoopPolicy` for Windows - Introduced a `get_event_loop()` utility to ensure proper event loop management, reusing a global loop on Windows and using process-specific loops on Linux, which improves stability and compatibility for asynchronous tasks across different environments. --- api.py | 128 +++++++++++++++++++------------ crawl.py | 4 +- job.py | 221 +++++++++++++++++++++++++++++++++--------------------- server.py | 55 ++++++++++---- tasks.py | 168 +++++++++++++++++------------------------ utils.py | 48 ++++++++---- 6 files changed, 357 insertions(+), 267 deletions(-) diff --git a/api.py b/api.py index 6173f9b..f0f7ce7 100644 --- a/api.py +++ b/api.py @@ -17,14 +17,15 @@ import os from typing import Any, AsyncGenerator, Dict, List, Optional, Tuple from urllib.parse import unquote -from uuid import uuid4 +from celery.result import AsyncResult # Import AsyncResult here +from celery_app import celery_app # Import celery_app here from courlan import get_base_url from crawl4ai import AsyncWebCrawler, BM25ContentFilter, BrowserConfig, CacheMode, CrawlResult, CrawlerRunConfig, DefaultMarkdownGenerator, LLMConfig, LLMContentFilter, LLMExtractionStrategy, LXMLWebScrapingStrategy, MemoryAdaptiveDispatcher, PruningContentFilter, RateLimiter from crawl4ai.utils import perform_completion_with_backoff from schemas import CrawlOperation from fastapi import HTTPException, Request,status from fastapi.responses import JSONResponse, StreamingResponse -import psutil +import signal import time from datetime import datetime # Import datetime for Celery task ID generation from celery.result import AsyncResult # Import AsyncResult @@ -504,11 +505,19 @@ async def cancel_a_job( if not celery_task.ready(): # Revoke the Celery task - import signal - celery_task.revoke( - terminate=True, - signal=signal.SIGKILL if force else signal.SIGTERM if os.name == 'posix' else signal.SIGTERM - ) + # Windows-compatible task revocation + if os.name == 'nt': # Windows + celery_task.revoke(terminate=True) + + # Force termination on Windows needs special handling + if force: + # Send shutdown event to worker + celery_app.control.broadcast('shutdown') + else: # POSIX systems + celery_task.revoke( + terminate=True, + signal=signal.SIGKILL if force else signal.SIGTERM + ) while True: # Check the task status to see if it's finished @@ -624,45 +633,72 @@ async def handle_stream_task_status( task_id, base_url: str = "", ): - - - task = decode_redis_hash(task) - async def stream_task_status(): - while True: - # Check the task status to see if it's finished - from celery.result import AsyncResult # Import AsyncResult here - from celery_app import celery_app # Import celery_app here - celery_task = AsyncResult(task_id, app=celery_app) # Re-initialize celery_task here - - if celery_task.ready(): - # You can get the final status and yield it - final_response = create_task_status_response(celery_task, task, task_id, base_url) - data = f"data: {json.dumps(final_response)}\n" - - yield data - break # Exit the loop when the task is complete - - # Yield the current status while waiting - response = create_task_status_response(celery_task, task, task_id, base_url) - data = f"data: {json.dumps(response)}\n" - - logger.info(f"send current status of celery app {data}") - - yield data - - await asyncio.sleep(1) # Wait before checking again - - yield "data: [DONE]" - - return StreamingResponse( - stream_task_status(), - media_type="text/event-stream", - headers={ - "Cache-Control": "no-cache", - "Connection": "keep-alive", - "X-Stream-Status": "active", - } - ) + """Stream status updates for a task.""" + try: + task = decode_redis_hash(task) + + async def stream_task_status(): + try: + while True: + # Check the task status to see if it's finished + celery_task = AsyncResult(task_id, app=celery_app) # Re-initialize celery_task here + + try: + celery_task_ready = celery_task.ready() + logger.info(f"Task {task_id}: Current status {celery_task.state}") + # Log status updates less frequently + if celery_task_ready: + logger.info(f"Task {task_id}: Completed with status {celery_task.state}") + + # Generate status response + response = create_task_status_response(celery_task, task, task_id, base_url) + data = f"data: {json.dumps(response)}\n" # Note the double newline for SSE format + + + yield data.encode('utf-8') # Ensure we're yielding bytes + + # Exit when the task is complete + if celery_task_ready: + break + + except Exception as e: + # Handle any serialization errors + error_msg = f"Error generating status response: {str(e)}" + logger.error(error_msg) + yield f"data: {json.dumps({'error': error_msg})}\n".encode('utf-8') + + # Wait before checking again + await asyncio.sleep(1) + + # Send the [DONE] marker to end the stream + yield b"data: [DONE]\n" + + except Exception as e: + # TODO: Handle exceptions in the streaming loop IN the frontend + logger.error(f"Fatal error in status stream: {str(e)}", exc_info=True) + yield f"event: error\ndata: {json.dumps({'error': str(e), 'fatal': True})}\n".encode('utf-8') + yield b"data: [DONE]\n" + + return StreamingResponse( + stream_task_status(), + media_type="text/event-stream", + headers={ + "Cache-Control": "no-cache, no-transform", + "Connection": "keep-alive", + "X-Accel-Buffering": "no", # Disable proxy buffering + "X-Stream-Status": "active", + } + ) + except Exception as e: + # Return a proper error response instead of letting the exception bubble up + logger.error(f"Error setting up status stream for task {task_id}: {str(e)}", exc_info=True) + return JSONResponse( + status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, + content={ + "status": "error", + "error": f"Failed to stream task status: {str(e)}" + } + ) async def create_new_task( redis: Redis, diff --git a/crawl.py b/crawl.py index f1d93bd..92ff73a 100644 --- a/crawl.py +++ b/crawl.py @@ -12,7 +12,7 @@ DefaultMarkdownGenerator, PruningContentFilter, ) -from fastapi import Request, Response +from fastapi import Request, Response, status from fastapi.responses import StreamingResponse import psutil from functools import partial @@ -132,7 +132,7 @@ async def reader(request: Request, response: Response) -> Response: # Return 406 Not Acceptable if Accept header is not supported return Response( - b"Unsupported Accept header", media_type="application/json", status_code=406 + b"Unsupported Accept header", media_type="application/json", status_code=status.HTTP_406_NOT_ACCEPTABLE ) except Exception as e: diff --git a/job.py b/job.py index 5a8d8f2..e96d624 100644 --- a/job.py +++ b/job.py @@ -3,7 +3,6 @@ from asyncio.log import logger import json import logging -import time from typing import Any, Callable, Dict, Optional, Union from celery import uuid @@ -12,7 +11,7 @@ from pydantic import BaseModel, HttpUrl from fastapi import APIRouter, Depends, HTTPException, Request, Response, WebSocket, WebSocketDisconnect, status -from redisCache import REDIS_CHANNEL, redis, pure_redis +from redisCache import REDIS_CHANNEL, redis, pure_redis, redis_xread from api import cancel_a_job, handle_crawl_job, handle_crawl_stream_job, handle_llm_request, handle_markdown_request, handle_stream_task_status, handle_task_status from auth import get_token_dependency from crawl import reader @@ -21,7 +20,9 @@ from triggers import event_stream from utils import load_config, safe_eval_config, setup_logging, stream_results -import gzip + +from celery.result import AsyncResult # Import AsyncResult here +from celery_app import celery_app # Import celery_app here config = load_config() setup_logging(config) @@ -70,7 +71,7 @@ def init_job_router(config, socket_client: set[Any]) -> APIRouter: verify_token = get_token_dependency(config) -# ---------- General endpoints ---------------------------------------------- +# ---------- General endpoints FOR Testing PURPOSE -------------------------- @job_router.get("/user/data") async def get_user_data( request: Request, @@ -288,7 +289,7 @@ async def get_task_id( # Cancel and Status general API ENDPOINTS -@job_router.put("/crawl/job/cancel/{temp_task_id}") +@job_router.put("/crawl/job/{temp_task_id}/cancel") async def crawl_job_cancel( request: Request, temp_task_id: str, @@ -319,7 +320,6 @@ async def crawl_job_cancel( } ) - # get status of a crawl job using celery task id @job_router.get("/crawl/job/status/{task_id}") @@ -338,30 +338,30 @@ async def crawl_stream_job_status( temp_task_id: str, decoded_token: Dict = Depends(verify_token) ): + """Get status of a crawl stream job using temporary task id with SSE.""" retries = 0 while retries < 3: try: # Get the task from redis task_id = await redis.hget(key=f"temp_task_id:{temp_task_id}", field='celery_task_id') - logger.info(task_id) + if not task_id: + logger.info(f"No task_id found for temp_task_id: {temp_task_id}") + retries += 1 + logger.info(f"Retrying to fetch task_id. Attempt {retries}/3") + await asyncio.sleep(1) # Wait before retrying + continue + + logger.info(f"Found task_id: {task_id} for temp_task_id: {temp_task_id}") task = await redis.hgetall(f"task:{task_id}") - if not task_id or task_id == 'empty' or not task: + if task_id == 'empty' or not task: retries += 1 - logger.info(f"Retrying to fetch task_id for temp_task_id: {temp_task_id}. Attempt {retries}/3") + logger.info(f"Task {task_id} not found or empty. Retrying {retries}/3") await asyncio.sleep(1) # Wait before retrying continue - if not task: - return JSONResponse( - status_code=status.HTTP_404_NOT_FOUND, - content={ - "status": "error", - "error": "Task not found" - } - ) - + # Successfully found task, return the streaming response return await handle_stream_task_status(task, task_id, base_url=str(request.base_url)) except Exception as e: @@ -370,11 +370,12 @@ async def crawl_stream_job_status( status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, content={ "status": "error", - "error": "Internal server error", + "error": "Internal server error while streaming task status", "internal_message": str(e) } ) + # After all retries, return not found return JSONResponse( status_code=status.HTTP_404_NOT_FOUND, content={ @@ -389,89 +390,137 @@ async def stream_crawl_results( task_id: str, decoded_token: bool = Depends(verify_token), ): - + channel = f"{REDIS_CHANNEL}:{task_id}" # Unique channel for the task async def event_stream(channel: str): + completed_yielded = False # Flag to indicate if the completion message has been yielded seen_messages = set() retries = 0 # Initialize retry counter - while True: - - await asyncio.sleep(1) # Simulate waiting for new messages - try: - - messages = await pure_redis.xread({channel: '0'}, count=None, block=5000) - - if retries > 12 and not completed_yielded: - logger.info("No completed message received after multiple retries. Ending stream.") - break - - logger.info(f"Received messages from Redis: {len(messages)} messages") - - if messages and isinstance(messages, list): - for message in messages: - if len(message) < 2: - logger.warning(f"Unexpected message format: {message}") - continue - - _, message_list = message - - for msg_id, msg_data in message_list: - if isinstance(msg_data, bytes): - msg_data_dict = json.loads(msg_data.decode("utf-8")) - elif isinstance(msg_data, dict): - msg_data_dict = msg_data - else: - logger.warning(f"Unexpected msg_data format: {msg_data}") - continue + last_id = "0" # Start from the beginning; use ">" for only new messages after consumer group creation + + + # Redis stream reading configuration + poll_interval = 0.5 # Reduced for lower latency + count = 20 # Increased for better throughput + block_ms = 5000 # 5 seconds block time + max_retries = 10 + """ For real-time streaming, count=15 is a good default: not too small, not too large. + If your stream is very high volume, you might increase it (e.g., count=100). + If you want lower latency (faster delivery per message), you might decrease it (e.g., count=1). """ + + # Check the task status to see if it's finished + celery_task = AsyncResult(task_id, app=celery_app) # Re-initialize celery_task here + try: + + while True: + + try: + # Yield heartbeat every 30 seconds to keep connection alive + # if retries > 0 and retries % 30 == 0: + # yield b"data: {\"type\":\"heartbeat\"}\n" + + # messages = await redis_xread(redis, {channel: '0'}, count=None, block=10000) + if retries > max_retries and not completed_yielded or (celery_task.ready() and retries > 3): + logger.info(f"Task {task_id}: Ending stream after {retries} retries with no activity") + # if not completed_yielded: + # yield b"data: {\"message\":\"completed\",\"type\":\"auto_complete\"}\n\n" + # yield b"data: [DONE]\n\n" + break + elif retries > max_retries and not completed_yielded and celery_task.state == "PENDING" or celery_task.state == "STARTED": + retries = 0 + + + # Read from Redis stream + messages = await pure_redis.xread({channel: last_id}, count, block=block_ms) + + + if messages and isinstance(messages, list): + # Reset retry counter when we get messages + retries = 0 + logger.info(f"Received messages from Redis: {len(messages)} messages") + + for _, message_list in messages: + for msg_id, msg_data in message_list: + last_id = msg_id # Update last_id after each message + + try: + if isinstance(msg_data, bytes): + msg_data_dict = json.loads(msg_data.decode("utf-8")) + elif isinstance(msg_data, dict): + msg_data_dict = msg_data + else: + logger.warning(f"Unexpected msg_data format: {msg_data}") + continue + except json.JSONDecodeError as e: + logger.error(f"JSON decode error: {e} for message: {msg_data}") + continue + + # Handle completion message + if msg_data_dict.get("message") == "completed": + if not completed_yielded: + logger.info(f"Task {task_id}: Yielding completion message") + # ADD 'data: ' PREFIX HERE + yield f"data: {json.dumps(msg_data_dict, ensure_ascii=False)}\n".encode('utf-8') + completed_yielded = True + yield b"data: [DONE]\n" + return + continue + + # Create unique ID to detect duplicates + unique_id = ( + f"{msg_data_dict.get('chunk_index', '')}_{msg_data_dict.get('url', msg_id)}" + if "chunk_index" in msg_data_dict + else msg_data_dict.get("id", msg_data_dict.get("url", msg_id)) + ) + + # Skip duplicates + if unique_id in seen_messages: + # logger.info(f"Duplicate message ignored {msg_id}: {unique_id}") + continue + + # Add to seen messages and yield to client + seen_messages.add(unique_id) + + # ADD 'data: ' PREFIX HERE + yield f"data: {json.dumps(msg_data_dict, ensure_ascii=False)}\n".encode('utf-8') + else: + # No messages - increment retry counter + retries += 1 + logger.warning(f"No messages returned or malformed response. Retry count: {retries}") - if isinstance(msg_data_dict, dict) and "message" in msg_data_dict and msg_data_dict["message"] == "completed": - if not completed_yielded: - logger.info("Yielding completed message.") - # ADD 'data: ' PREFIX HERE - yield ("data: " + json.dumps(msg_data_dict, ensure_ascii=False) + "\n").encode('utf-8') - completed_yielded = True - break - continue - - unique_id = f"{msg_data_dict.get('chunk_index', '')}_{msg_data_dict.get('url', msg_id)}" \ - if "chunk_index" in msg_data_dict else msg_data_dict.get("id",msg_data_dict.get("url", msg_id)) - - if unique_id in seen_messages: - logger.info(f"Duplicate message ignored {msg_id}: {unique_id}") - continue - - seen_messages.add(unique_id) - - logger.info(f"Received message on str {msg_id}") - - retries = 0 - # ADD 'data: ' PREFIX HERE - yield ("data: " + json.dumps(msg_data_dict, ensure_ascii=False) + "\n").encode('utf-8') - else: - logger.warning("No messages returned or malformed response.") - - except Exception as e: - logger.error(f"Error in event stream: {e}") - - finally: - logger.info("Closing the event stream.") + except Exception as e: + logger.error(f"Error reading from Redis stream: {e}") + # Send error to client as an SSE event + yield f"event: error\ndata: {json.dumps({'error': str(e)})}\n".encode('utf-8') + # Optionally end the stream after a serious error + retries += 1 + + # Pause before next poll + await asyncio.sleep(poll_interval) + + except asyncio.CancelledError: + logger.info(f"Task {task_id}: Stream cancelled by client") + yield b"data: {\"message\":\"stream_cancelled\"}\n" - if completed_yielded : + except Exception as e: + logger.error(f"Fatal error in event stream: {e}", exc_info=True) + yield f"event: error\ndata: {json.dumps({'error': str(e), 'fatal': True})}\n".encode('utf-8') + finally: + seen_messages.clear() + # Always send DONE if we exit the loop without returning + if not completed_yielded: yield b"data: [DONE]\n" - break - else: - retries += 1 - logger.info(f"No completed message received. Continuing to listen for new messages. Retry count: {retries}") + logger.info(f"Task {task_id}: Stream closed") - seen_messages.clear() return StreamingResponse( event_stream(channel), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", + "X-Accel-Buffering": "no", # Disable proxy buffering "X-Stream-Status": "active", }) diff --git a/server.py b/server.py index e6746e5..938ccee 100644 --- a/server.py +++ b/server.py @@ -8,7 +8,6 @@ import sys import time from typing import Any, Callable -import signal # from typing import Annotated # noqa: F401 from crawl4ai import AsyncWebCrawler, BrowserConfig @@ -17,20 +16,12 @@ FastAPI, HTTPException, Request, + WebSocket, status ) from fastapi.responses import RedirectResponse, JSONResponse from fastapi.exceptions import RequestValidationError from starlette.datastructures import Address # Import Address -from fastapi import ( - Depends, - FastAPI, - HTTPException, - Request, - status -) -from fastapi.responses import RedirectResponse, JSONResponse -from fastapi.exceptions import RequestValidationError from fastapi.middleware.cors import CORSMiddleware from fastapi.middleware.gzip import GZipMiddleware from fastapi.middleware.httpsredirect import HTTPSRedirectMiddleware @@ -70,7 +61,7 @@ logger = logging.getLogger(__name__) -__version__ = config["app"]["version"] or "0.5.1-d1" +__version__ = config["app"]["version"] or "0.0.1-d1" # ── global page semaphore (hard cap) ───────────────────────── MAX_PAGES = config["crawler"]["pool"].get("max_pages", 30) @@ -85,11 +76,19 @@ async def capped_arun(self, *a, **kw): AsyncWebCrawler.arun = capped_arun +orig_arun_many = AsyncWebCrawler.arun_many + +async def capped_arun_many(self, urls, config=None, dispatcher=None, **kwargs): + async with GLOBAL_SEM: + return await orig_arun_many(self, urls, config, dispatcher, **kwargs) +AsyncWebCrawler.arun_many = capped_arun_many + + # Set the number of workers NUM_WORKERS = int(os.getenv("NUM_WORKERS", os.cpu_count() or 1)) # Store connected WebSocket clients -socket_client = set() +socket_client: set[WebSocket] = set() if sys.platform != "win32": import uvloop # type: ignore @@ -99,6 +98,12 @@ async def capped_arun(self, *a, **kw): asyncio.set_event_loop_policy(EventLoopPolicy()) # logger.warning("uvloop is not supported on Windows, using default(auto) event loop") +production = os.getenv("PYTHON_ENV", "development").lower() == "production" +if production: + print("\033[94mINFO-SERVER:\033[0m \033[92mRunning in Production mode. PYTHON_ENV\033[0m", production) +else: + print("\033[94mWARNIN-SERVER:\033[92m Running in Development mode. PYTHON_ENV", production) + ############################################################### # ───────────────────── FastAPI lifespan ────────────────────── ############################################################### @@ -113,7 +118,7 @@ async def lifespan(_: FastAPI): **config["crawler"]["browser"].get("kwargs", {}), )) # warm‑up await test_connection(redis) # Moved from on_event("startup") - await test_connection(pure_redis) # Moved from on_event("startup") + # await test_connection(pure_redis) # Moved from on_event("startup") app.state.janitor = asyncio.create_task(janitor()) # idle GC app.state.websocket = asyncio.create_task(periodic_client_cleanup(socket_client)) yield @@ -274,6 +279,17 @@ def _setup_security(app_: FastAPI): # ───────────────────── FastAPI middlewares ────────────────────── ################################################################ +@app.middleware("http") +async def add_process_time_header(request: Request, call_next): + + start_time = time.time() + response = await call_next(request) + process_time = time.time() - start_time + response.headers["X-Process-Time"] = f"{process_time * 1000:.2f}" + + print(f"Request: {request.url.path} - Response time: {process_time * 1000:.2f} ms") + return response + # security headers middleware @app.middleware("http") async def add_security_headers(request: Request, call_next): @@ -334,7 +350,7 @@ async def rate_limit_middleware(request: Request, call_next): # Add CORS middleware app.add_middleware( CORSMiddleware, - allow_origins=config["app"].get("cors_origins", ["*"]), + allow_origins=config["app"].get(("cors_origins" if production else "cors_origins_dev"), ["*"]), allow_credentials=True, allow_methods=["*"], allow_headers=["*"], @@ -365,8 +381,15 @@ async def root(decoded_token: bool = Depends(verify_token)): # health check endpoint @app.get(config["observability"]["health_check"]["endpoint"]) -async def health(): - return {"status": "ok", "timestamp": time.time(), "version": __version__} +async def health(request: Request, response: JSONResponse): + """Health check endpoint.""" + try: + return JSONResponse({"status": "ok", "timestamp": time.time(), "version": __version__}) + + except Exception as e: + logger.error(f"Health check failed: {e}", exc_info=True) + response.status_code = status.HTTP_503_SERVICE_UNAVAILABLE + return JSONResponse({"status": "error", "detail": str(e)}) # prometheus metrics endpoint @app.get(config["observability"]["prometheus"]["endpoint"]) diff --git a/tasks.py b/tasks.py index b055261..01deeb4 100644 --- a/tasks.py +++ b/tasks.py @@ -4,6 +4,7 @@ import logging import os from functools import partial # Import partial +import platform import sys import time from typing import AsyncGenerator, Dict, List, Optional, cast @@ -27,6 +28,49 @@ from utils import FilterType, TaskStatus, _get_memory_mb, datetime_handler, decode_redis_hash, is_task_id, setup_logging, should_cleanup_task, load_config, stream_pubsub_results, task_status_color from crawler_pool import get_crawler, cancel_crawler +if sys.platform != "win32": + import uvloop # type: ignore + asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) +else: + from asyncio import WindowsProactorEventLoopPolicy as EventLoopPolicy + asyncio.set_event_loop_policy(EventLoopPolicy()) + +# Track if we're on Windows +IS_WINDOWS = platform.system() == "Windows" + +# Only use a global event loop on Windows with pool=solo +# On Linux with concurrency, each worker process will manage its own loop +_event_loop = None + +def get_event_loop(): + """ + Get or create an event loop in a platform-specific way. + + On Windows with pool=solo: Returns a persistent global event loop + On Linux with multiprocessing: Returns a process-specific event loop + """ + global _event_loop + + if IS_WINDOWS: + # Windows approach: reuse the same event loop for all tasks + if _event_loop is None or _event_loop.is_closed(): + _event_loop = asyncio.new_event_loop() + asyncio.set_event_loop(_event_loop) + return _event_loop + else: + # Linux approach: get the current process's event loop + try: + loop = asyncio.get_event_loop() + if loop.is_closed(): + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + return loop + except RuntimeError: + # No event loop in current thread, create one + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + return loop + # # Global Redis clients for all Celery tasks # redis_url = os.environ.get("UPSTASH_REDIS_REST_URL") # redis_token = os.environ.get("UPSTASH_REDIS_REST_TOKEN") @@ -46,7 +90,7 @@ # allow_telemetry=False, # ) -# # Global PureRedis client +# # # Global PureRedis client # pure_redis = PureRedis( # host=str(REDIS_URL), # port=int(REDIS_PORT), @@ -78,12 +122,7 @@ def get_redis(): redis, pure_redis = get_redis() -if sys.platform == "win32": - asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()) -else: - import uvloop # type: ignore - asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) - + async def get_firebase_client(): global _firebase_client, _db_instance if _firebase_client is None: @@ -91,7 +130,7 @@ async def get_firebase_client(): _db_instance, _ = _firebase_client.init_firebase() return _db_instance -@celery_app.task(bind=True) +@celery_app.task(bind=True, ) def crawl_task(self, urls: List[str], browser_config: Dict, crawler_config: Dict) : """Celery task to handle crawl requests.""" logger.info("Starting crawl task with URLs") @@ -104,16 +143,6 @@ def crawl_task(self, urls: List[str], browser_config: Dict, crawler_config: Dict # # await redis.hset(f"task:{task_id}", values={"status": TaskStatus.CANCELED}) # return {"status": TaskStatus.CANCELED, "message": "Task aborted by user."} - - # return {"status": TaskStatus.COMPLETED, "message": "Crawl task completed successfully."} - # async def runner(): - # import time - # data = "No URLs provided" - # for i in range(10): # Long-running task - # time.sleep(1) - # print(f"Crawling {data}, step {i}") - # return f"Crawled {data}" - asyncio.run(_crawl_task_impl(self, urls, browser_config, crawler_config)) return {"status": TaskStatus.COMPLETED, "message": "Crawl task completed successfully."} @@ -193,9 +222,25 @@ def crawl_stream_task(self, uid: str, operation_id: str, urls: List[str], browse """Celery task to process crawl stream.""" try: - # Use the global event loop set at module level - loop = asyncio.get_event_loop() - loop.run_until_complete(_crawl_stream_task_impl(self, uid, operation_id, urls, browser_config, crawler_config)) + # Get an appropriate event loop for the platform + loop = get_event_loop() + + # Run the implementation + loop.run_until_complete( + _crawl_stream_task_impl(self, uid, operation_id, urls, browser_config, crawler_config) + ) + + # On Linux, we need to clean up pending tasks before returning + # On Windows with solo pool, we keep tasks running + if not IS_WINDOWS: + pending = [task for task in asyncio.all_tasks(loop) + if not task.done() and task is not asyncio.current_task(loop)] + if pending: + logger.info(f"Cleaning up {len(pending)} pending tasks") + for task in pending: + task.cancel() + loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True)) + return {"status": TaskStatus.COMPLETED, "message": "crawl stream task completed successfully."} except Exception as e: @@ -462,7 +507,6 @@ async def _crawl_stream_task_impl( await monitor.record_metrics(operResult) await cancel_crawler(sign) # Remove the crawler to free resources - await asyncio.sleep(5) # Give Redis time to process the update # FIXME: Update operation status CHECK THE RESULTS @@ -533,84 +577,6 @@ async def _crawl_stream_task_impl( await cancel_crawler(sign) # Remove the crawler to free resources raise -""" async def handle_stream_crawl_request( - urls: List[str], - crawler: AsyncWebCrawler, - crawler_config: dict, - config: dict -) -> tuple[AsyncGenerator, float, Optional[float], Optional[float]]: - Handle non-streaming crawl requests - start_mem_mb = _get_memory_mb() # <--- Get memory before - start_time = time.time() - mem_delta_mb = None - peak_mem_mb = start_mem_mb - try: - urls = [('https://' + url) if not url.startswith(('http://', 'https://')) else url for url in urls] - - crawler_conf = CrawlerRunConfig.load(crawler_config) - - dispatcher = MemoryAdaptiveDispatcher( - memory_threshold_percent=config["crawler"]["memory_threshold_percent"], - rate_limiter=RateLimiter( - base_delay=tuple(config["crawler"]["rate_limiter"]["base_delay"]) - ) if config["crawler"]["rate_limiter"]["enabled"] else None - ) - - - base_config = config["crawler"]["base_config"] - # Iterate on key-value pairs in global_config then use haseattr to set them - for key, value in base_config.items(): - if hasattr(crawler_conf, key): - setattr(crawler_conf, key, value) - - - func = getattr(crawler, "arun" if len(urls) == 1 else "arun_many") - partial_func = partial(func, - urls[0] if len(urls) == 1 else urls, - config=crawler_conf, - dispatcher=dispatcher) - - results = await partial_func() - - # await crawler.close() - - end_mem_mb = _get_memory_mb() # <--- Get memory after - end_time = time.time() - total_time = end_time - start_time - - if start_mem_mb is not None and end_mem_mb is not None: - mem_delta_mb = end_mem_mb - start_mem_mb # <--- Calculate delta - peak_mem_mb = max(peak_mem_mb if peak_mem_mb else 0, end_mem_mb) # <--- Get peak memory - logger.info(f"Memory usage: Start: {start_mem_mb} MB, End: {end_mem_mb} MB, Delta: {mem_delta_mb} MB, Peak: {peak_mem_mb} MB, Total Time: {total_time}" ) - - # return [{"url": url, - # "dump": result.model_dump() if hasattr(result, 'model_dump') else "No model_dump available", - # } for url, result in zip(urls, results)], total_time, mem_delta_mb, peak_mem_mb - return results, total_time, mem_delta_mb, peak_mem_mb - - except Exception as e: - logger.error(f"Crawl error: {str(e)}", exc_info=True) - if 'crawler' in locals() and crawler.ready: # Check if crawler was initialized and started - try: - await crawler.close() - except Exception as e: - logger.error(f"Error closing crawler during exception handling: {str(e)}") - - # Measure memory even on error if possible - end_mem_mb_error = _get_memory_mb() - if start_mem_mb is not None and end_mem_mb_error is not None: - mem_delta_mb = end_mem_mb_error - start_mem_mb - - raise HTTPException( - status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, - detail=json.dumps({ # Send structured error - "error": str(e), - "server_memory_delta_mb": mem_delta_mb, - "server_peak_memory_mb": max(peak_mem_mb if peak_mem_mb else 0, end_mem_mb_error or 0) - }) - ) """ - - async def handle_stream_crawl_request( urls: List[str], crawler:AsyncWebCrawler, @@ -662,7 +628,7 @@ async def handle_stream_crawl_request( # Raising HTTPException here will prevent streaming response raise HTTPException( status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, - detail=json.dumps({ # Send structured error + detail=json.dumps({ # Send structured error "error": str(e), "server_memory_delta_mb": mem_delta_mb, # "server_peak_memory_mb": max(peak_mem_mb if peak_mem_mb else 0, end_mem_mb_error or 0) diff --git a/utils.py b/utils.py index 2159bd1..e32e7da 100644 --- a/utils.py +++ b/utils.py @@ -5,6 +5,7 @@ import logging import os import re +from fastapi import WebSocket import psutil # from upstash_redis.asyncio import Redis from redis.asyncio import Redis # Use redis.asyncio for async Redis operations @@ -17,6 +18,8 @@ import urllib.parse import random +from redisCache import redis_xadd + logger = logging.getLogger(__name__) class TaskStatus(str, Enum): @@ -87,10 +90,10 @@ def setup_logging(config: Dict) -> None: level=config["logging"]["level"], format=config["logging"]["format"] ) -async def remove_stale_clients(socket_client) -> None: +async def remove_stale_clients(socket_client: set[WebSocket]) -> None: """Remove stale WebSocket clients.""" - disconnected_clients = set() + disconnected_clients:set[WebSocket] = set() for client in socket_client: try: await client.send_text("ping") # Ping the client @@ -237,7 +240,7 @@ def create_task_status_response(celery_task:AsyncResult, task: Dict[str, str], t "task_id": task_id, "status": convert_celery_status(celery_task.state) or task["status"], "created_at": task["created_at"], - "urls": task["urls"], + "urls": task.get("urls", ""), "_links": { "self": {"href": f"{base_url}llm/{task_id}"}, "refresh": {"href": f"{base_url}llm/{task_id}"} @@ -245,14 +248,23 @@ def create_task_status_response(celery_task:AsyncResult, task: Dict[str, str], t } if task["status"] == TaskStatus.COMPLETED or celery_task.ready(): - response["result"] = celery_task.result if celery_task.successful() else None or json.loads(task["result"]) + # Safely parse result if it exists and isn't empty + if celery_task.successful() and celery_task.result: + response["result"] = celery_task.result + elif task.get("result") and task["result"].strip(): + try: + response["result"] = json.loads(task["result"]) + except json.JSONDecodeError: + # If we can't parse the result as JSON, use it as a string + response["result"] = task["result"] + else: + response["result"] = None elif task["status"] == TaskStatus.FAILED or celery_task.failed(): - response["error"] = task["error"] + response["error"] = task.get("error", "Unknown error") response["result"] = celery_task.result return response - async def stream_results(crawler: _c4.AsyncWebCrawler, results_gen: AsyncGenerator) -> AsyncGenerator[bytes, None]: """Stream results with heartbeats and completion markers.""" import json @@ -297,8 +309,8 @@ async def stream_pubsub_results(redis: Redis, channel: str, results_gen: AsyncGe result: _c4.CrawlResult complete = {"status": "ok", "message": "completed"} - # data: list[dict[str, Any]] = [] - buffer: list[dict[str, Any]] = [] + # buffer: list[dict[str, Any]] = [] + pipe2 = redis.pipeline() try: async for result in results_gen: try: @@ -309,10 +321,9 @@ async def stream_pubsub_results(redis: Redis, channel: str, results_gen: AsyncGe # remove html from result before sending to redis result_dict["html"] = "" # type: ignore - # result_dict['server_memory_mb'] = server_memory_mb + result_dict['server_memory_mb'] = server_memory_mb # result_dict['status'] = "model_dump" url = result_dict.get('url', 'unknown') - logger.info(f"Publishing result for {url}") model_dump = result_dict if hasattr(result, 'model_dump') \ else {"status": "error", "message": "No model_dump available skipping", @@ -320,12 +331,14 @@ async def stream_pubsub_results(redis: Redis, channel: str, results_gen: AsyncGe if isinstance(model_dump, dict): # data.append(model_dump) - buffer.append(model_dump) + logger.info(f"Publishing result for {url}") + # buffer.append(model_dump) else: raise ValueError(model_dump) batch_json = json.dumps(model_dump, default=datetime_handler, ensure_ascii=False) pipe = redis.pipeline() + total_chunks = (len(batch_json) + chunk_size - 1) // chunk_size # Calculate total chunks # Split batch_json into chunks of chunk_size for i in range(0, len(batch_json), chunk_size): chunk = batch_json[i:i+chunk_size] @@ -335,27 +348,30 @@ async def stream_pubsub_results(redis: Redis, channel: str, results_gen: AsyncGe "type": "batch_chunk", "url": url, "chunk_index": str(i // chunk_size), + "total_chunks": str(total_chunks), # Add total_chunks attribute "dump": chunk #.encode("utf-8") if isinstance(chunk, str) else chunk }) await pipe.execute() - buffer.clear() except Exception as e: logger.error(f"Serialization error: {e}") error_response = {"status": "error", "message": str(e), "url": getattr(result, 'url', 'unknown')} - await redis.xadd(channel, {key: str(value) if isinstance(value, bool) else value for key, value in error_response.items() } ) + + pipe2.xadd(channel, {key: str(value) if isinstance(value, bool) else value for key, value in error_response.items() } ) complete = {"status": "error", "message": "completed"} - await redis.xadd(channel, {key: str(value) if isinstance(value, bool) else value for key, value in complete.items()}) + pipe2.xadd(channel, {key: str(value) if isinstance(value, bool) else value for key, value in complete.items()}) except asyncio.CancelledError: logger.warning("Client disconnected during streaming") - await redis.xadd(channel, {"status": "canceled", "message": "streaming canceled"}) + pipe2.xadd(channel, {"status": "canceled", "message": "streaming canceled"}) except Exception as e: logger.error(f"Unexpected error in stream_pubsub_results: {e}") - await redis.xadd(channel, {"status": "error", "message": str(e)}) + pipe2.xadd(channel, {"status": "error", "message": str(e)}) + await pipe2.execute() return False + await pipe2.execute() return True From 3f6c2217e54feb004ca1a952e2ec716f0b4b2528 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Mon, 22 Sep 2025 16:06:12 +0300 Subject: [PATCH 06/30] feat(celery): enhance environment variable loading and add Windows-specific signal handling --- celery_app.py | 52 ++++++++++++++++++++++++++++++++++++++++++--------- 1 file changed, 43 insertions(+), 9 deletions(-) diff --git a/celery_app.py b/celery_app.py index af45ce2..d94bbdd 100644 --- a/celery_app.py +++ b/celery_app.py @@ -13,9 +13,11 @@ # else: # import uvloop # type: ignore # asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) - -# ────────────────── configuration ────────────────── -load_dotenv(verbose=True) +production = os.getenv("PYTHON_ENV", "development").lower() == "production" +env_file = ".env" if production else "dev.env" + +# Load environment variables +load_dotenv(env_file, verbose=True) # Celery needs a URI for broker and backend, even if we're passing an instance. # For Upstash Redis, the URL is typically what's needed. @@ -35,11 +37,12 @@ REDIS_URI = f"rediss://{REDIS_USERNAME}:{REDIS_PASSWORD}@{REDIS_URL}:{REDIS_PORT}/0?ssl_cert_reqs=CERT_REQUIRED" celery_app = Celery( - "deepcrawl4ai", + "crawlagent", broker=REDIS_URI, backend=REDIS_URI, - include=["tasks"] # We will create a tasks.py later + include=["tasks"], # We will create a tasks.py later ) +# celery_app.config_from_object('celeryconfig') celery_app.conf.update( task_track_started=True, @@ -48,13 +51,43 @@ task_serializer='json', result_serializer='json', accept_content=['json'], + # broker_pool_limit = 2, # Default is 10 + # broker_connection_max_retries = 3, # Default is 100 + # broker_heartbeat = 0, # Disabled by default + # broker_transport_options = {'visibility_timeout': 3600}, # 1 hour timezone='UTC', enable_utc=True, - # task_soft_time_limit=600, # Soft time limit: 10 minutes - # task_time_limit=900, # Hard time limit: 15 minutes - # result_expires=3600, # Expire results after 1 hour to avoid memory bloat + task_soft_time_limit=600, # Soft time limit: 10 minutes + task_time_limit=900, # Hard time limit: 15 minutes + result_expires=3600, # Expire results after 1 hour to avoid memory bloat + # broker_connection_retry=True, + # broker_connection_retry_on_startup=True, + # broker_connection_max_retries=10, + # broker_connection_timeout=30, + # broker_transport_options={ + # 'visibility_timeout': 3600, + # 'socket_timeout': 30, + # 'socket_connect_timeout': 30, + # }, + # redis_socket_keepalive=True, + + # Windows-specific settings + worker_cancel_long_running_tasks_on_connection_loss=True, # Helps with Windows task cancellation + task_remote_tracebacks=True, # Better error reporting + worker_max_tasks_per_child=1 if os.name == 'nt' else None, # Prevent memory leaks on Windows ) +# Add Windows-specific signal handling +if os.name == 'nt': + # Windows doesn't support SIGKILL/SIGTERM the same way + def windows_shutdown_handler(*args): + print("Windows shutdown signal received") + celery_app.control.broadcast('shutdown') + sys.exit(0) + + signal.signal(signal.SIGTERM, windows_shutdown_handler) + signal.signal(signal.SIGINT, windows_shutdown_handler) + def graceful_shutdown(signum, frame): logging.info(f"Received signal {signum}, shutting down Celery worker gracefully...") from celery.worker import state @@ -63,7 +96,8 @@ def graceful_shutdown(signum, frame): signal.signal(signal.SIGTERM, graceful_shutdown) signal.signal(signal.SIGINT, graceful_shutdown) - + + # celery_app.conf.task_queues = ( # Queue("light_jobs"), # Queue("heavy_jobs"), From 04f36b33fc490c9c19dcb01aa6eba5195b49617c Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Mon, 22 Sep 2025 16:06:49 +0300 Subject: [PATCH 07/30] feat(docker): upgrade base image to Python 3.12 and optimize Dockerfile structure - Upgraded the base Python image from `3.10-slim` to `3.12-slim` for both build and final stages, leveraging the latest Python features and performance. - Consolidated system dependency installations into a single `RUN` layer to optimize image caching and reduce layers. - Added essential build dependencies (`build-essential`, `cmake`, `gcc`, `g++`, etc.) for potential C extensions and core utilities (`lsof`, `ca-certificates`). - Streamlined the `pip wheel` and installation process, removing redundant steps and cleanup commands. - Added `DISPLAY=:99` environment variable for headless browser environments. --- Dockerfile | 239 ++++++++++++++++++++++++------------------------- Dockerfile.old | 179 ++++++++++++++++++++++++++++++++++++ 2 files changed, 296 insertions(+), 122 deletions(-) create mode 100644 Dockerfile.old diff --git a/Dockerfile b/Dockerfile index a42fdef..a238f52 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,7 +1,7 @@ -# Use an official Python runtime as a base image -FROM python:3.10-slim AS build +# Build stage +FROM python:3.12-slim AS builder -# Set environment variables to production +# Set build-time environment variables ENV PYTHONFAULTHANDLER=1 \ PYTHONHASHSEED=random \ PYTHONUNBUFFERED=1 \ @@ -9,142 +9,137 @@ ENV PYTHONFAULTHANDLER=1 \ PYTHONDONTWRITEBYTECODE=1 \ PIP_DISABLE_PIP_VERSION_CHECK=1 \ PIP_DEFAULT_TIMEOUT=100 \ - DEBIAN_FRONTEND=noninteractive \ + DEBIAN_FRONTEND=noninteractive + +# Set build arguments +ARG APP_HOME=/app + +# Create app directory +WORKDIR ${APP_HOME} + +# Install build dependencies +COPY requirements.txt . +RUN pip wheel --no-cache-dir --no-deps --wheel-dir /app/wheels -r requirements.txt + +# Final stage +FROM python:3.12-slim + +# Metadata +LABEL maintainer="Prokopis Antoniadis" \ + description="🔥🕷️ Crawl4AI: LLM Web Crawler & scraper" \ + version="1.0" + +# Set environment variables +ENV PYTHONUNBUFFERED=1 \ + PYTHONDONTWRITEBYTECODE=1 \ PLAYWRIGHT_BROWSERS_PATH=/ms-playwright \ - PYTHON_ENV=production + PYTHON_ENV=production \ + DISPLAY=:99 +# Set build arguments ARG APP_HOME=/app -ARG PYTHON_VERSION=3.12 -ARG INSTALL_TYPE=default -ARG ENABLE_GPU=false -ARG TARGETARCH - - -# Add Maintainer Info -LABEL maintainer="Prokopis Antoniadis" -LABEL description="🔥🕷️ Crawl4AI: LLM Web Crawler & scraper" -LABEL version="1.0" - -# RUN apt-get update && apt-get install -y --no-install-recommends \ -# build-essential \ -# curl \ -# wget \ -# gnupg \ -# cmake \ -# pkg-config \ -# python3-dev \ -# libjpeg-dev \ -# supervisor \ -# && apt-get clean \ -# && rm -rf /var/lib/apt/lists/* - -# Install minimal system dependencies for Playwright and Python -RUN apt-get update -y && \ - apt-get upgrade -y && \ - apt-get install -y --no-install-recommends \ + +WORKDIR ${APP_HOME} + +# Install system dependencies in a single layer +RUN apt-get update && apt-get install -y --no-install-recommends \ + fonts-liberation \ + ca-certificates \ + lsof \ + # Add build dependencies for madoka + build-essential \ + curl \ + wget \ + gnupg \ + cmake \ + gcc \ + g++ \ + pkg-config \ + python3-dev \ + libjpeg-dev \ + supervisor \ libnss3 \ + # Playwright system dependencies + libglib2.0-0 \ + libnss3 \ + libnspr4 \ + libatk1.0-0 \ + libatk-bridge2.0-0 \ + libcups2 \ + libdrm2 \ + libdbus-1-3 \ + libxcb1 \ + libxkbcommon0 \ + libx11-6 \ libxcomposite1 \ - libxcursor1 \ libxdamage1 \ + libxext6 \ + libxfixes3 \ libxrandr2 \ - libdrm2 \ libgbm1 \ - libxss1 \ + libpango-1.0-0 \ + libcairo2 \ libasound2 \ - libatk1.0-0 \ - libatk-bridge2.0-0 \ + libatspi2.0-0 \ + libxcursor1 \ + libxss1 \ libgtk-3-0 \ xvfb \ x11vnc \ git \ - curl fonts-liberation ca-certificates lsof \ - && git clone https://github.com/novnc/noVNC /opt/noVNC \ - && git clone https://github.com/novnc/websockify /opt/noVNC/utils/websockify \ - && apt-get clean && rm -rf /var/lib/apt/lists/* - -# RUN apt-get update && apt-get dist-upgrade -y \ -# && rm -rf /var/lib/apt/lists/* - -# RUN if [ "$ENABLE_GPU" = "true" ] && [ "$TARGETARCH" = "amd64" ] ; then \ -# apt-get update && apt-get install -y --no-install-recommends \ -# nvidia-cuda-toolkit \ -# && apt-get clean \ -# && rm -rf /var/lib/apt/lists/* ; \ -# else \ -# echo "Skipping NVIDIA CUDA Toolkit installation (unsupported platform or GPU disabled)"; \ -# fi - -# RUN if [ "$TARGETARCH" = "arm64" ]; then \ -# echo "🦾 Installing ARM-specific optimizations"; \ -# apt-get update && apt-get install -y --no-install-recommends \ -# libopenblas-dev \ -# && apt-get clean \ -# && rm -rf /var/lib/apt/lists/*; \ -# elif [ "$TARGETARCH" = "amd64" ]; then \ -# echo "🖥️ Installing AMD64-specific optimizations"; \ -# apt-get update && apt-get install -y --no-install-recommends \ -# libomp-dev \ -# && apt-get clean \ -# && rm -rf /var/lib/apt/lists/*; \ -# else \ -# echo "Skipping platform-specific optimizations (unsupported platform)"; \ -# fi - -# Create a non-root user and group -# RUN groupadd -r appuser && useradd --no-log-init -r -g appuser appuser - -# Create and set permissions for appuser home directory -# RUN mkdir -p /home/appuser && chown -R appuser:appuser /home/appuser - -# Set the working directory inside the container -WORKDIR ${APP_HOME} - -# Copy the requirements file and install Python dependencies -# Copy supervisor config first (might need root later, but okay for now) -# COPY supervisord.conf . - -COPY requirements.txt . -COPY config.yml . - -RUN pip install --no-cache-dir -r requirements.txt && \ - pip install --no-cache-dir playwright - -RUN pip install --no-cache-dir --upgrade pip && \ - python -c "import crawl4ai; print('✅ crawl4ai is ready to rock!')" && \ - python -c "from playwright.sync_api import sync_playwright; print('✅ Playwright is feeling dramatic!')" - - -# RUN crawl4ai-setup - -# https://playwright.dev/docs/browsers -# Install only the required Playwright browser (Chromium) -RUN playwright install chromium - -# RUN playwright install --with-deps - + # Add sudo for X11 management + && git clone --depth 1 https://github.com/novnc/noVNC /opt/noVNC \ + && git clone --depth 1 https://github.com/novnc/websockify /opt/noVNC/utils/websockify \ + && apt-get clean \ + && rm -rf /var/lib/apt/lists/* \ + && rm -rf /var/cache/apt/* + +# Create non-root user +RUN groupadd -r appuser && \ + useradd --no-log-init -r -g appuser appuser && \ + mkdir -p /home/appuser/.cache /ms-playwright && \ + chown -R appuser:appuser /home/appuser /ms-playwright ${APP_HOME} + +# Copy wheels from builder +COPY --from=builder /app/wheels /wheels +COPY --from=builder /app/requirements.txt . + +# Install dependencies +RUN pip install --no-cache-dir /wheels/* \ + && rm -rf /wheels \ + && pip install --no-cache-dir playwright \ + && playwright install --with-deps chromium \ + && python -c "import crawl4ai; print('✅ crawl4ai is ready to rock!')" \ + && python -c "from playwright.sync_api import sync_playwright; print('✅ Playwright is feeling dramatic!')" + +# Copy application code +COPY --chown=appuser:appuser . . +COPY --chown=appuser:appuser config.yml . + +# Set display environment variable +ENV DISPLAY=:99 + +# Run diagnostics RUN crawl4ai-doctor -# RUN mkdir -p /home/appuser/.cache/ms-playwright \ -# && cp -r /root/.cache/ms-playwright/chromium-* /home/appuser/.cache/ms-playwright/ \ -# && chown -R appuser:appuser /home/appuser/.cache/ms-playwright +# Expose ports +EXPOSE 8000 9222 6080 -# Copy the application code into the container -COPY . . -# Change ownership of the application directory to the non-root user -# RUN chown -R appuser:appuser ${APP_HOME} +# Healthcheck dont need, fly io do this for us +# HEALTHCHECK --interval=30s --timeout=30s --start-period=5s --retries=3 \ +# CMD curl -f http://localhost:8000/health || exit 1 -# Switch to the non-root user before starting the application -# USER appuser +# Copy and set permissions for the entrypoint script +COPY --chown=root:root docker-entrypoint.sh /usr/local/bin/docker-entrypoint.sh +RUN chmod +x /usr/local/bin/docker-entrypoint.sh -# # Install only the required Playwright browser (Chromium) -# RUN playwright install chromium +# Switch to non-root user +USER appuser -# Expose the port your FastAPI app will run on -EXPOSE 8000 9222 6080 +# Set the entrypoint +ENTRYPOINT ["docker-entrypoint.sh"] -# Command to run the application -# CMD ["uvicorn", "server:app", "--host", "0.0.0.0", "--port", "8000", "--ws", "websockets"] -# CMD ["sh", "-c", "Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] -CMD ["sh", "-c", "Xvfb :99 -screen 0 1280x720x24 & x11vnc -display :99 -nopw -forever -shared -rfbport 5900 -quiet & /opt/noVNC/utils/websockify/run 6080 localhost:5900 --web /opt/noVNC & uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] -# Start the application using supervisord -# CMD ["supervisord", "-c", "supervisord.conf"] +# Start application +CMD ["sh", "-c", "\ + /opt/noVNC/utils/websockify/run --web /opt/noVNC 0.0.0.0:6080 0.0.0.0:5900 & \ + uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] diff --git a/Dockerfile.old b/Dockerfile.old new file mode 100644 index 0000000..921bfff --- /dev/null +++ b/Dockerfile.old @@ -0,0 +1,179 @@ +# Use an official Python runtime as a base image +FROM python:3.10-slim AS build + +# Set environment variables to production +# Set build-time environment variables +ENV PYTHONFAULTHANDLER=1 \ + PYTHONHASHSEED=random \ + PYTHONUNBUFFERED=1 \ + PIP_NO_CACHE_DIR=1 \ + PYTHONDONTWRITEBYTECODE=1 \ + PIP_DISABLE_PIP_VERSION_CHECK=1 \ + PIP_DEFAULT_TIMEOUT=100 \ + DEBIAN_FRONTEND=noninteractive \ + PLAYWRIGHT_BROWSERS_PATH=/ms-playwright \ + PYTHON_ENV=production \ + DISPLAY=:99 + +ARG APP_HOME=/app +ARG PYTHON_VERSION=3.10 +ARG INSTALL_TYPE=default +ARG ENABLE_GPU=false +ARG TARGETARCH + + +# Add Maintainer Info +LABEL maintainer="Prokopis Antoniadis" +LABEL description="🔥🕷️ Crawl4AI: LLM Web Crawler & scraper" +LABEL version="1.0" + +RUN apt-get update && apt-get install -y --no-install-recommends \ + fonts-liberation \ + ca-certificates \ + lsof \ + # Add build dependencies for madoka + build-essential \ + curl \ + wget \ + gnupg \ + cmake \ + gcc \ + g++ \ + pkg-config \ + python3-dev \ + libjpeg-dev \ + supervisor \ + && apt-get clean \ + && rm -rf /var/lib/apt/lists/* + +# Install minimal system dependencies for Playwright and Python +RUN apt-get update -y && \ + apt-get upgrade -y && \ + apt-get install -y --no-install-recommends \ + libnss3 \ + libxcomposite1 \ + libxcursor1 \ + libxdamage1 \ + libxrandr2 \ + libdrm2 \ + libgbm1 \ + libxss1 \ + libasound2 \ + libatk1.0-0 \ + libatk-bridge2.0-0 \ + libgtk-3-0 \ + xvfb \ + x11vnc \ + git \ + curl fonts-liberation ca-certificates lsof \ + && git clone https://github.com/novnc/noVNC /opt/noVNC \ + && git clone https://github.com/novnc/websockify /opt/noVNC/utils/websockify \ + && apt-get clean && rm -rf /var/lib/apt/lists/* + +# RUN apt-get update && apt-get dist-upgrade -y \ +# && rm -rf /var/lib/apt/lists/* + +# RUN if [ "$ENABLE_GPU" = "true" ] && [ "$TARGETARCH" = "amd64" ] ; then \ +# apt-get update && apt-get install -y --no-install-recommends \ +# nvidia-cuda-toolkit \ +# && apt-get clean \ +# && rm -rf /var/lib/apt/lists/* ; \ +# else \ +# echo "Skipping NVIDIA CUDA Toolkit installation (unsupported platform or GPU disabled)"; \ +# fi + +# RUN if [ "$TARGETARCH" = "arm64" ]; then \ +# echo "🦾 Installing ARM-specific optimizations"; \ +# apt-get update && apt-get install -y --no-install-recommends \ +# libopenblas-dev \ +# && apt-get clean \ +# && rm -rf /var/lib/apt/lists/*; \ +# elif [ "$TARGETARCH" = "amd64" ]; then \ +# echo "🖥️ Installing AMD64-specific optimizations"; \ +# apt-get update && apt-get install -y --no-install-recommends \ +# libomp-dev \ +# && apt-get clean \ +# && rm -rf /var/lib/apt/lists/*; \ +# else \ +# echo "Skipping platform-specific optimizations (unsupported platform)"; \ +# fi + +# Create a non-root user and group +RUN groupadd -r appuser && useradd --no-log-init -r -g appuser appuser + +# Create and set permissions for appuser home directory +RUN mkdir -p /home/appuser && chown -R appuser:appuser /home/appuser + +# Set the working directory inside the container +WORKDIR ${APP_HOME} + +# Copy the requirements file and install Python dependencies +# Copy supervisor config first (might need root later, but okay for now) +# COPY supervisord.conf . + +COPY requirements.txt . +COPY config.yml . + +RUN pip install --no-cache-dir -r requirements.txt && \ + pip install --no-cache-dir playwright + +RUN pip install --no-cache-dir --upgrade pip && \ + python -c "import crawl4ai; print('✅ crawl4ai is ready to rock!')" && \ + python -c "from playwright.sync_api import sync_playwright; print('✅ Playwright is feeling dramatic!')" + + +# RUN crawl4ai-setup + +# https://playwright.dev/docs/browsers +# Install only the required Playwright browser (Chromium) +# Install and verify Playwright +RUN playwright install --with-deps && \ + python -c "from playwright.sync_api import sync_playwright; \ + with sync_playwright() as p: \ + browser = p.chromium.launch(); \ + browser.close(); \ + print('✅ Playwright browser verification successful!')" + +# Create necessary directories and set permissions before switching to non-root user +RUN mkdir -p /home/appuser/.cache /ms-playwright && \ + chown -R appuser:appuser /home/appuser /ms-playwright /app + +RUN crawl4ai-doctor + +# RUN mkdir -p /home/appuser/.cache/ms-playwright \ +# && cp -r /root/.cache/ms-playwright/chromium-* /home/appuser/.cache/ms-playwright/ \ +# && chown -R appuser:appuser /home/appuser/.cache/ms-playwright + +# Change ownership of the application directory to the non-root user +RUN chown -R appuser:appuser ${APP_HOME} + +# Copy the application code into the container +COPY --chown=appuser:appuser . . + +# Expose the port your FastAPI app will run on +EXPOSE 8000 9222 6080 + +# Healthcheck +HEALTHCHECK --interval=30s --timeout=30s --start-period=5s --retries=3 \ + CMD bash -c '\ + MEM=$(free -m | awk "/^Mem:/{print \$2}"); \ + if [ $MEM -lt 2048 ]; then \ + echo "⚠️ Warning: Less than 2GB RAM available! Your container might need a memory boost! 🚀"; \ + exit 1; \ + fi && \ + curl -f http://localhost:8000/health || exit 1' + + +# Switch to the non-root user before starting the application +USER appuser + +# Set environment variables to ptoduction +ENV PYTHON_ENV=production + +# Command to run the application +# CMD ["uvicorn", "server:app", "--host", "0.0.0.0", "--port", "8000", "--ws", "websockets"] +# CMD ["sh", "-c", "Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] +# gunicorn server:app -w 4 -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000 +CMD ["sh", "-c", "Xvfb :99 -screen 0 1280x720x24 & x11vnc -display :99 -nopw -forever -shared -rfbport 5900 -quiet & /opt/noVNC/utils/websockify/run 0.0.0.0:6080 0.0.0.0:5900 --web /opt/noVNC & uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] +# Start the application using supervisord +# CMD ["supervisord", "-c", "supervisord.conf"] From 1e207408bdf86fec9e26dc1eb8612061b50aa654 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 02:44:51 +0300 Subject: [PATCH 08/30] =?UTF-8?q?The=20variable=C2=A0chunk=5Fsize=C2=A0is?= =?UTF-8?q?=20used=20but=20not=20defined=20in=20this=20function.=20It=20sh?= =?UTF-8?q?ould=20be=20declared=20as=20a=20parameter=20or=20defined=20as?= =?UTF-8?q?=20a=20constant=20within=20the=20function=20scope=20to=20avoid?= =?UTF-8?q?=20potential=20NameError.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- utils.py | 1 + 1 file changed, 1 insertion(+) diff --git a/utils.py b/utils.py index e32e7da..4f36fec 100644 --- a/utils.py +++ b/utils.py @@ -338,6 +338,7 @@ async def stream_pubsub_results(redis: Redis, channel: str, results_gen: AsyncGe batch_json = json.dumps(model_dump, default=datetime_handler, ensure_ascii=False) pipe = redis.pipeline() + chunk_size = 4096 # Define chunk_size as a constant (adjust as needed) total_chunks = (len(batch_json) + chunk_size - 1) // chunk_size # Calculate total chunks # Split batch_json into chunks of chunk_size for i in range(0, len(batch_json), chunk_size): From 568b22414f28eaad2fe370f6de3bf53de21a9c79 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 02:48:07 +0300 Subject: [PATCH 09/30] =?UTF-8?q?The=C2=A0response=C2=A0parameter=20is=20d?= =?UTF-8?q?eclared=20but=20not=20used=20properly.=20The=20function=20creat?= =?UTF-8?q?es=20a=20new=C2=A0JSONResponse=C2=A0instead=20of=20modifying=20?= =?UTF-8?q?the=20passed=20response=20parameter.=20Either=20remove=20the=20?= =?UTF-8?q?parameter=20or=20use=20it=20consistently=20throughout=20the=20f?= =?UTF-8?q?unction.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- server.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/server.py b/server.py index 938ccee..6bd10eb 100644 --- a/server.py +++ b/server.py @@ -381,15 +381,15 @@ async def root(decoded_token: bool = Depends(verify_token)): # health check endpoint @app.get(config["observability"]["health_check"]["endpoint"]) -async def health(request: Request, response: JSONResponse): +async def health(request: Request): """Health check endpoint.""" try: return JSONResponse({"status": "ok", "timestamp": time.time(), "version": __version__}) except Exception as e: logger.error(f"Health check failed: {e}", exc_info=True) - response.status_code = status.HTTP_503_SERVICE_UNAVAILABLE - return JSONResponse({"status": "error", "detail": str(e)}) + # status code set in JSONResponse below + return JSONResponse({"status": "error", "detail": str(e)}, status_code=status.HTTP_503_SERVICE_UNAVAILABLE) # prometheus metrics endpoint @app.get(config["observability"]["prometheus"]["endpoint"]) From c735661fc6b65de686665fca97ff0bb3d05ab173 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 02:52:01 +0300 Subject: [PATCH 10/30] =?UTF-8?q?The=20parentheses=20around=20the=20ternar?= =?UTF-8?q?y=20expression=20create=20a=20tuple=20instead=20of=20selecting?= =?UTF-8?q?=20the=20correct=20key.=20It=20should=20be=C2=A0config[\"app\"]?= =?UTF-8?q?.get(\"cors=5Forigins\"=20if=20production=20else=20\"cors=5Fori?= =?UTF-8?q?gins=5Fdev\",=20[\"*\"]).?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- server.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server.py b/server.py index 6bd10eb..7001c22 100644 --- a/server.py +++ b/server.py @@ -350,7 +350,7 @@ async def rate_limit_middleware(request: Request, call_next): # Add CORS middleware app.add_middleware( CORSMiddleware, - allow_origins=config["app"].get(("cors_origins" if production else "cors_origins_dev"), ["*"]), + allow_origins=config["app"].get("cors_origins" if production else "cors_origins_dev", ["*"]), allow_credentials=True, allow_methods=["*"], allow_headers=["*"], From 9abf3ca588441a097d34c09cd71ad17552138968 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 02:53:10 +0300 Subject: [PATCH 11/30] =?UTF-8?q?The=20boolean=20logic=20is=20unclear=20du?= =?UTF-8?q?e=20to=20operator=20precedence.=20The=20condition=20should=20us?= =?UTF-8?q?e=20parentheses=20to=20clarify=20the=20intended=20logic:=C2=A0i?= =?UTF-8?q?f=20(retries=20>=20max=5Fretries=20and=20not=20completed=5Fyiel?= =?UTF-8?q?ded)=20or=20(celery=5Ftask.ready()=20and=20retries=20>=203):.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- job.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/job.py b/job.py index e96d624..f5c131c 100644 --- a/job.py +++ b/job.py @@ -422,7 +422,7 @@ async def event_stream(channel: str): # yield b"data: {\"type\":\"heartbeat\"}\n" # messages = await redis_xread(redis, {channel: '0'}, count=None, block=10000) - if retries > max_retries and not completed_yielded or (celery_task.ready() and retries > 3): + if (retries > max_retries and not completed_yielded) or (celery_task.ready() and retries > 3): logger.info(f"Task {task_id}: Ending stream after {retries} retries with no activity") # if not completed_yielded: # yield b"data: {\"message\":\"completed\",\"type\":\"auto_complete\"}\n\n" From a598e3710b716b036bc0945440f89a3685ac857e Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 21:25:19 +0300 Subject: [PATCH 12/30] Potential fix for code scanning alert no. 7: Information exposure through an exception Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com> --- api.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/api.py b/api.py index f0f7ce7..4b427d3 100644 --- a/api.py +++ b/api.py @@ -665,7 +665,7 @@ async def stream_task_status(): # Handle any serialization errors error_msg = f"Error generating status response: {str(e)}" logger.error(error_msg) - yield f"data: {json.dumps({'error': error_msg})}\n".encode('utf-8') + yield f"data: {json.dumps({'error': 'An internal error occurred while generating the status response.'})}\n".encode('utf-8') # Wait before checking again await asyncio.sleep(1) @@ -676,7 +676,7 @@ async def stream_task_status(): except Exception as e: # TODO: Handle exceptions in the streaming loop IN the frontend logger.error(f"Fatal error in status stream: {str(e)}", exc_info=True) - yield f"event: error\ndata: {json.dumps({'error': str(e), 'fatal': True})}\n".encode('utf-8') + yield f"event: error\ndata: {json.dumps({'error': 'A fatal error occurred while streaming the task status.', 'fatal': True})}\n".encode('utf-8') yield b"data: [DONE]\n" return StreamingResponse( From 8ab13c0456c2033f3db9cd4dffaab0d238347521 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 21:38:24 +0300 Subject: [PATCH 13/30] XADD path: missing await and wrong API surface. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit pipe.execute([...]) won’t work for either client. Call the right method and await it. Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- redisCache.py | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/redisCache.py b/redisCache.py index 11c9a13..b4262db 100644 --- a/redisCache.py +++ b/redisCache.py @@ -148,7 +148,18 @@ async def redis_xadd(pipe, channel: str, message: dict, maxlen: Optional[int] = for key, value in message.items(): pieces.extend([key, value]) - message_id = pipe.execute(["XADD", *pieces]) + # Upstash client exposes `execute`; redis-py uses `execute_command` or `xadd` + if hasattr(pipe, "execute"): + message_id = await pipe.execute("XADD", *pieces) + elif hasattr(pipe, "xadd"): + kwargs = {} + if maxlen is not None: + kwargs["maxlen"] = maxlen + kwargs["approximate"] = approximate + message_id = await pipe.xadd(channel, message, **kwargs) + else: + # Fallback to redis-py low-level + message_id = await pipe.execute_command("XADD", *pieces) return message_id except Exception as e: print(f"\033[91mERROR-DB:\033[0m Failed to add message to stream '{channel}': {e}") From f9bf6badcca960060cf6304033d5c7830990e664 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 21:39:34 +0300 Subject: [PATCH 14/30] XREAD path: flatten args and support both clients. Passing the whole list to execute will fail. Build varargs and branch per client. Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- redisCache.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/redisCache.py b/redisCache.py index b4262db..ba1e4c3 100644 --- a/redisCache.py +++ b/redisCache.py @@ -166,7 +166,7 @@ async def redis_xadd(pipe, channel: str, message: dict, maxlen: Optional[int] = return None -async def redis_xread(redis: Redis, streams: dict, count: Optional[int] = None, block: Optional[int] = None): +async def redis_xread(redis: Redis | PureRedis, streams: dict, count: Optional[int] = None, block: Optional[int] = None): """ Read messages from a Redis stream using XREAD. :param redis: Redis client (PureRedis) @@ -194,7 +194,11 @@ async def redis_xread(redis: Redis, streams: dict, count: Optional[int] = None, command.append(channel) for channel, last_id in streams.items(): command.append(last_id) - result = await redis.execute(command) + # Flatten into varargs for the underlying client + if hasattr(redis, "execute"): + result = await redis.execute(command[0], *command[1:]) + else: + result = await redis.execute_command(*command) return result except Exception as e: print(f"\033[91mERROR-DB:\033[0m Failed to read from stream(s) '{list(streams.keys())}': {e}") From 24f49e9667fbd4b51cbf2524e627f9451d920b47 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 21:40:20 +0300 Subject: [PATCH 15/30] fix(server): Health endpoint: mask internal errors and silence ARG001. Returning str(e) leaks details (CodeQL finding); also request is unused. Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- server.py | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/server.py b/server.py index 7001c22..68dfe9b 100644 --- a/server.py +++ b/server.py @@ -381,15 +381,16 @@ async def root(decoded_token: bool = Depends(verify_token)): # health check endpoint @app.get(config["observability"]["health_check"]["endpoint"]) -async def health(request: Request): +@app.get(config["observability"]["health_check"]["endpoint"]) +async def health(_: Request): """Health check endpoint.""" try: return JSONResponse({"status": "ok", "timestamp": time.time(), "version": __version__}) - except Exception as e: - logger.error(f"Health check failed: {e}", exc_info=True) - # status code set in JSONResponse below - return JSONResponse({"status": "error", "detail": str(e)}, status_code=status.HTTP_503_SERVICE_UNAVAILABLE) + except Exception: + logger.exception("Health check failed") + # Do not expose internal error details to clients + return JSONResponse({"status": "error"}, status_code=status.HTTP_503_SERVICE_UNAVAILABLE) # prometheus metrics endpoint @app.get(config["observability"]["prometheus"]["endpoint"]) From 8a87c330229c997f666dca46b74f73cea4a030b0 Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 21:41:14 +0300 Subject: [PATCH 16/30] fix(dockerfile): Builder stage lacks build deps; wheels may fail to build. Packages needing compilation (e.g., with C extensions) will fail without toolchains/headers. Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- Dockerfile | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/Dockerfile b/Dockerfile index a238f52..7a6ce9f 100644 --- a/Dockerfile +++ b/Dockerfile @@ -19,8 +19,10 @@ WORKDIR ${APP_HOME} # Install build dependencies COPY requirements.txt . -RUN pip wheel --no-cache-dir --no-deps --wheel-dir /app/wheels -r requirements.txt - +RUN apt-get update && apt-get install -y --no-install-recommends \ + build-essential gcc g++ python3-dev pkg-config libjpeg-dev cmake \ + && rm -rf /var/lib/apt/lists/* \ + && pip wheel --no-cache-dir --no-deps --wheel-dir /app/wheels -r requirements.txt # Final stage FROM python:3.12-slim From 9b156e80611d6693c291dcc3de59e661241a1e7f Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 21:46:56 +0300 Subject: [PATCH 17/30] Potential fix for code scanning alert no. 8: Information exposure through an exception Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com> --- api.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/api.py b/api.py index 4b427d3..63e7dd9 100644 --- a/api.py +++ b/api.py @@ -696,7 +696,7 @@ async def stream_task_status(): status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, content={ "status": "error", - "error": f"Failed to stream task status: {str(e)}" + "error": "Failed to stream task status due to an internal error." } ) From 6e7b28561e44267d1f21f8fba0ecc0ff2f8709fd Mon Sep 17 00:00:00 2001 From: Prokopis Antoniadis Date: Tue, 23 Sep 2025 22:14:42 +0300 Subject: [PATCH 18/30] =?UTF-8?q?fix(redis):=20XADD=20path:=20missing=20re?= =?UTF-8?q?quired=20=E2=80=9C*=E2=80=9D=20ID=20and=20wrong=20method=20pref?= =?UTF-8?q?erence.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Raw XADD requires an ID; you’re not sending “*”. Prefer high‑level xadd when available; only fall back to raw execute. Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> --- redisCache.py | 17 ++++++++++------- 1 file changed, 10 insertions(+), 7 deletions(-) diff --git a/redisCache.py b/redisCache.py index ba1e4c3..e6d56f6 100644 --- a/redisCache.py +++ b/redisCache.py @@ -145,21 +145,24 @@ async def redis_xadd(pipe, channel: str, message: dict, maxlen: Optional[int] = if approximate: pieces.append("~") pieces.append(str(maxlen)) + # When not specifying an explicit ID, Redis requires "*" + pieces.append("*") for key, value in message.items(): - pieces.extend([key, value]) + pieces.extend([str(key), str(value)]) - # Upstash client exposes `execute`; redis-py uses `execute_command` or `xadd` - if hasattr(pipe, "execute"): - message_id = await pipe.execute("XADD", *pieces) - elif hasattr(pipe, "xadd"): + # Prefer high-level `xadd` when available + if hasattr(pipe, "xadd"): kwargs = {} if maxlen is not None: kwargs["maxlen"] = maxlen kwargs["approximate"] = approximate message_id = await pipe.xadd(channel, message, **kwargs) else: - # Fallback to redis-py low-level - message_id = await pipe.execute_command("XADD", *pieces) + # Fallback: Upstash `execute` or redis-py `execute_command` + if hasattr(pipe, "execute"): + message_id = await pipe.execute("XADD", *pieces) + else: + message_id = await pipe.execute_command("XADD", *pieces) return message_id except Exception as e: print(f"\033[91mERROR-DB:\033[0m Failed to add message to stream '{channel}': {e}") From cbe9a7e582fffb74cc2eb9e21045ea08b7a8f9ce Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:06:39 +0300 Subject: [PATCH 19/30] fix(docker-entrypoint): add process management and signal handling - Implement proper background process management with graceful shutdown and VNC authentication support --- docker-entrypoint.sh | 73 ++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 71 insertions(+), 2 deletions(-) diff --git a/docker-entrypoint.sh b/docker-entrypoint.sh index 4d1cf56..3c3b754 100644 --- a/docker-entrypoint.sh +++ b/docker-entrypoint.sh @@ -1,4 +1,8 @@ #!/bin/sh +set -eu + +# Initialize NOVNC_PID at the top of the script +NOVNC_PID="" # Ensure the X11 socket directory exists and has correct permissions # This is done as root before switching to appuser @@ -6,9 +10,74 @@ mkdir -p /tmp/.X11-unix chmod 1777 /tmp/.X11-unix chown appuser:appuser /tmp/.X11-unix -# Start Xvfb and x11vnc in the background +# Start Xvfb in the background Xvfb :99 -screen 0 1280x720x24 -ac & -x11vnc -display :99 -nopw -forever -shared -rfbport 5900 -quiet & +XVFB_PID=$! + +# Setup VNC authentication +# Optional: set VNC_PASSWORD via env; default to disabling if missing in production +if [ -n "${VNC_PASSWORD:-}" ]; then + mkdir -p "$HOME/.vnc" + x11vnc -storepasswd "$VNC_PASSWORD" "$HOME/.vnc/passwd" >/dev/null 2>&1 + AUTH_ARGS="-rfbauth $HOME/.vnc/passwd" +else + # In dev only; never leave unauthenticated in production + AUTH_ARGS="-nopw" +fi + +# Start x11vnc with auth in the background +# Bind to localhost; expose only via websockify +x11vnc -display :99 ${AUTH_ARGS} -forever -shared -rfbport 5900 -localhost -quiet & +X11VNC_PID=$! + +# Define cleanup function +cleanup() { + echo "Shutting down background processes..." + + # Kill processes with proper signal handling + if [ -n "$XVFB_PID" ]; then + kill -TERM $XVFB_PID 2>/dev/null || true + fi + + if [ -n "$X11VNC_PID" ]; then + kill -TERM $X11VNC_PID 2>/dev/null || true + fi + + if [ -n "$NOVNC_PID" ]; then + kill -TERM $NOVNC_PID 2>/dev/null || true + fi + + # Wait for processes to terminate + wait $XVFB_PID 2>/dev/null || true + wait $X11VNC_PID 2>/dev/null || true + [ -n "$NOVNC_PID" ] && wait $NOVNC_PID 2>/dev/null || true + + echo "Cleanup complete." +} + +# Set up signal trapping +trap cleanup INT TERM QUIT + +# Export DISPLAY for child processes +export DISPLAY=:99 + +# Log successful setup +echo "✅ X11 environment setup complete. Display: $DISPLAY" + +# Check if this is a worker process +if echo "$@" | grep -q "celery"; then + echo "Running as worker process, starting minimal X11 services" + # Skip starting noVNC for worker processes +else + echo "Running as app process, starting all services" + # Start noVNC only for app process + if [ -d "/opt/noVNC" ]; then + /opt/noVNC/utils/websockify/run --web /opt/noVNC 0.0.0.0:6080 localhost:5900 & + NOVNC_PID=$! + else + NOVNC_PID="" + fi +fi # Execute the main command passed to the entrypoint exec "$@" From 0a642a40cbcc20c59666fd5cf05ec0df57750935 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:07:12 +0300 Subject: [PATCH 20/30] feat(storage): implement S3 multipart upload and streaming capabilities - Add robust S3 storage with multipart uploads, connection pooling, retry logic, and streaming decompression --- storage.py | 351 +++++++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 284 insertions(+), 67 deletions(-) diff --git a/storage.py b/storage.py index 63e9145..2808ff2 100644 --- a/storage.py +++ b/storage.py @@ -1,10 +1,15 @@ -from datetime import datetime +from datetime import datetime, timezone import os -from typing import Any, Dict +from typing import Any, Dict, Optional, List import aioboto3 import zstandard as zstd import pydantic +import logging +import asyncio +from botocore.exceptions import ClientError +# Set up logging +logger = logging.getLogger(__name__) class TigrisBucketResult(pydantic.BaseModel): key_name: str @@ -13,7 +18,8 @@ class TigrisBucketResult(pydantic.BaseModel): ge=0 ) # ensures file size is non-negative in kilobytes file_compressed_size: float = pydantic.Field(ge=0) - created_at: datetime = pydantic.Field(default_factory=datetime.now) + created_at: datetime = pydantic.Field(default_factory=lambda: datetime.now(timezone.utc)) + updated_at: Optional[datetime] = pydantic.Field(default_factory=lambda: datetime.now(timezone.utc)) def to_str(self) -> str: """Return a formatted string representation""" @@ -32,92 +38,303 @@ def to_dict(self) -> Dict[str, Any]: """Return dictionary representation""" return self.model_dump() - -# from botocore.exceptions import DataNotFoundError, UnknownServiceError - -# Tigris Buckets S3-compatible configuration -TIGRIS_BUCKET_NAME = "deepcrawl4.bucket.a" +# Tigris Buckets configuration +TIGRIS_BUCKET_NAME = "crawlagent.bucket.a" TIGRIS_ACCESS_KEY = os.getenv("AWS_ACCESS_KEY_ID") TIGRIS_SECRET_KEY = os.getenv("AWS_SECRET_ACCESS_KEY") -TIGRIS_ENDPOINT_URL = os.getenv("AWS_ENDPOINT_URL_S3") # Change if needed - +TIGRIS_ENDPOINT_URL = os.getenv("AWS_ENDPOINT_URL_S3") -# Initialize S3 client for Tigris +# Create a reusable session session = aioboto3.Session() +_client = None - -async def folder_exists(folder_name: str) -> bool: - """Check if a folder exists in Tigris Buckets.""" - try: - async with session.client( # type: ignore +# Get or create S3 client with proper context management +async def get_s3_client(): + """Get S3 client as a context manager.""" + global _client + if _client is None: + _client = await session.client( "s3", endpoint_url=TIGRIS_ENDPOINT_URL, aws_access_key_id=TIGRIS_ACCESS_KEY, aws_secret_access_key=TIGRIS_SECRET_KEY, - ) as svc: - # List buckets - response = await svc.list_buckets() - for bucket in response["Buckets"]: - print(f' {bucket["Name"]}') + ).__aenter__() + return _client + +# Use this function to create a context manager +class S3ClientManager: + async def __aenter__(self): + self.client = await get_s3_client() + return self.client + + async def __aexit__(self, exc_type, exc_val, exc_tb): + # Don't close the global client + pass - # List objects +async def folder_exists(folder_name: str) -> bool: + """Check if a folder exists in Tigris Buckets.""" + try: + async with S3ClientManager() as svc: response = await svc.list_objects_v2( - Bucket=TIGRIS_BUCKET_NAME, Prefix=f"{folder_name}/", MaxKeys=1 + Bucket=TIGRIS_BUCKET_NAME, + Prefix=f"{folder_name}/", + MaxKeys=1 ) - for obj in response["Contents"]: - print(f' {obj["Key"]}') - - return "Contents" in response # Returns True if folder has at least one file + return "Contents" in response + except ClientError as e: + logger.error(f"S3 error checking folder {folder_name}: {e}") + return False except Exception as e: - print(f"Error checking folder existence: {e}") + logger.error(f"Unexpected error checking folder {folder_name}: {e}") return False - -async def upload_markdown( +async def upload_compressed_file( markdown_str: str, folder_name: str, file_name: str -) -> TigrisBucketResult | None: - """Uploads a Markdown string directly to Tigris Buckets if it's >= 100 KB.""" - - file_size_kb = len(markdown_str.encode("utf-8")) / 1024 # Convert bytes to KB - - # if file_size_kb < 100: - # print( - # f"⚠️ Skipping upload: {file_name} is only {file_size_kb:.2f} KB (less than 100 KB)" - # ) - # return # Skip upload if file size is less than 100 KB - - key_name = f"{folder_name}/{file_name}" # E.g., "markdown-files/report.md" +) -> Optional[TigrisBucketResult]: + """Uploads a Markdown string compressed with zstd to Tigris Buckets.""" + + file_size_kb = len(markdown_str.encode("utf-8")) / 1024 + key_name = f"{folder_name}/{file_name}" + + # Compress data + try: + compressor = zstd.ZstdCompressor(level=3) # Level 3 balances speed and compression + compressed_data = compressor.compress(markdown_str.encode("utf-8")) + file_compressed_size = len(compressed_data) / 1024 + + # Determine if we should use multipart upload (>5MB) + use_multipart = len(compressed_data) > 5 * 1024 * 1024 + + async with S3ClientManager() as svc: + if use_multipart: + # For large files, use multipart upload + return await _upload_multipart( + svc, compressed_data, key_name, file_name, file_size_kb, file_compressed_size + ) + else: + # For smaller files, use simple put_object + await svc.put_object( + Bucket=TIGRIS_BUCKET_NAME, + Key=key_name, + Body=compressed_data, + ContentType="text/markdown", + ) + logger.info(f"Uploaded {file_name} to {key_name} ({file_size_kb:.2f} KB → {file_compressed_size:.2f} KB)") + + return TigrisBucketResult( + key_name=key_name, + file_name=file_name, + file_size=file_size_kb, + file_compressed_size=file_compressed_size, + ) + except ClientError as e: + logger.error(f"S3 upload error for {key_name}: {e}") + return None + except Exception as e: + logger.error(f"Compression/upload error for {file_name}: {e}") + return None +async def _upload_multipart(svc, data, key_name, file_name, file_size_kb, file_compressed_size): + """Helper function for multipart uploads of large files.""" try: - async with session.client( # type: ignore - "s3", - endpoint_url=TIGRIS_ENDPOINT_URL, - aws_access_key_id=TIGRIS_ACCESS_KEY, - aws_secret_access_key=TIGRIS_SECRET_KEY, - ) as svc: - # Zstd (Balanced for Speed & Compression) - # Pros: Faster and better compression than Gzip. - # Cons: Less widely supported than Gzip. - compressor = zstd.ZstdCompressor() - compressed_data = compressor.compress(markdown_str.encode("utf-8")) - file_compressed_size = len(compressed_data) / 1024 # Convert bytes to KB - print(f"Compressed {file_name} to {file_compressed_size:.2f} KB") - - await svc.put_object( + # Create multipart upload + mpu = await svc.create_multipart_upload( + Bucket=TIGRIS_BUCKET_NAME, + Key=key_name, + ContentType="text/markdown" + ) + + # Split data into chunks (5MB per chunk) + chunk_size = 5 * 1024 * 1024 + chunks = [data[i:i+chunk_size] for i in range(0, len(data), chunk_size)] + + # Upload parts + parts = [] + for i, chunk in enumerate(chunks): + part_number = i + 1 + response = await svc.upload_part( Bucket=TIGRIS_BUCKET_NAME, Key=key_name, - Body=compressed_data, - ContentType="text/markdown", + PartNumber=part_number, + UploadId=mpu["UploadId"], + Body=chunk ) - print( - f"✅ Uploaded {file_name} to {TIGRIS_BUCKET_NAME}/{key_name} ({file_size_kb:.2f} KB)" + parts.append({ + "PartNumber": part_number, + "ETag": response["ETag"] + }) + + # Complete multipart upload + await svc.complete_multipart_upload( + Bucket=TIGRIS_BUCKET_NAME, + Key=key_name, + UploadId=mpu["UploadId"], + MultipartUpload={"Parts": parts} + ) + + logger.info(f"Multipart uploaded {file_name} to {key_name} ({file_size_kb:.2f} KB → {file_compressed_size:.2f} KB)") + + return TigrisBucketResult( + key_name=key_name, + file_name=file_name, + file_size=file_size_kb, + file_compressed_size=file_compressed_size, + ) + except Exception as e: + logger.error(f"Multipart upload error: {e}") + # Try to abort the multipart upload to avoid orphaned uploads + try: + await svc.abort_multipart_upload( + Bucket=TIGRIS_BUCKET_NAME, + Key=key_name, + UploadId=mpu["UploadId"] ) - return TigrisBucketResult( - key_name=key_name, - file_name=file_name, - file_size=file_size_kb, - file_compressed_size=file_compressed_size, + except Exception as abort_error: + logger.error(f"Failed to abort multipart upload: {abort_error}") + return None + +async def download_file_decompressed(folder_name: str, file_name: str) -> Optional[str]: + """Downloads and decompresses a file from Tigris Buckets.""" + key_name = f"{folder_name}/{file_name}" + + # Implement exponential backoff retry for downloads + max_retries = 3 + retry_delay = 1 # Start with 1 second delay + + for attempt in range(max_retries): + try: + async with S3ClientManager() as svc: + response = await svc.get_object( + Bucket=TIGRIS_BUCKET_NAME, + Key=key_name + ) + compressed_data = await response["Body"].read() + + # Decompress the data + decompressor = zstd.ZstdDecompressor() + decompressed_data = decompressor.decompress(compressed_data) + str_file = decompressed_data.decode("utf-8") + + logger.info(f"Downloaded and decompressed {file_name} from {key_name}") + return str_file + + except ClientError as e: + if e.response['Error']['Code'] == 'NoSuchKey': + logger.error(f"File not found: {key_name}") + return None + logger.warning(f"S3 error on attempt {attempt+1}/{max_retries}: {e}") + except Exception as e: + logger.warning(f"Download error on attempt {attempt+1}/{max_retries}: {e}") + + # Only sleep if we're going to retry + if attempt < max_retries - 1: + await asyncio.sleep(retry_delay) + retry_delay *= 2 # Exponential backoff + + logger.error(f"Download failed after {max_retries} attempts: {key_name}") + return None + +async def list_files(folder_name: str) -> List[Dict[str, Any]]: + """List all files in a folder.""" + try: + async with await get_s3_client() as svc: + response = await svc.list_objects_v2( + Bucket=TIGRIS_BUCKET_NAME, + Prefix=f"{folder_name}/" ) + + if "Contents" not in response: + return [] + + return [ + { + "key": obj["Key"], + "size": obj["Size"], + "last_modified": obj["LastModified"], + "file_name": obj["Key"].split("/")[-1] + } + for obj in response["Contents"] + ] except Exception as e: - print(f"❌ Upload failed: {e}") + logger.error(f"Error listing files in {folder_name}: {e}") + return [] + +async def get_presigned_url(folder_name: str, file_name: str, expiration: int = 3600) -> Optional[str]: + """ + Generate a presigned URL for direct download of the compressed file. + + Args: + folder_name: The folder/prefix containing the file + file_name: The file name to download + expiration: URL expiration time in seconds (default 1 hour) + + Returns: + Presigned URL string or None if error + """ + key_name = f"{folder_name}/{file_name}" + + try: + async with S3ClientManager() as svc: + # Create the presigned URL + url = await svc.generate_presigned_url( + 'get_object', + Params={ + 'Bucket': TIGRIS_BUCKET_NAME, + 'Key': key_name + }, + ExpiresIn=expiration + ) + + logger.info(f"Generated presigned URL for {key_name}, expires in {expiration} seconds") + return url + + except Exception as e: + logger.error(f"Error generating presigned URL for {key_name}: {e}") return None + +async def download_and_decompress_stream(folder_name: str, file_name: str): + """ + Downloads and decompresses a file, returning it as a streaming response. + This function should be used with FastAPI's StreamingResponse. + + Args: + folder_name: The folder/prefix containing the file + file_name: The file name to download + + Returns: + An async generator yielding decompressed content + """ + key_name = f"{folder_name}/{file_name}" + + async def content_stream(): + try: + async with S3ClientManager() as svc: + response = await svc.get_object( + Bucket=TIGRIS_BUCKET_NAME, + Key=key_name + ) + + # Use zstd streaming decompression for memory efficiency + decompressor = zstd.ZstdDecompressor() + compressed_stream = response["Body"] + + # Read and decompress in chunks + chunk_size = 1024 * 1024 # 1MB chunks + while True: + chunk = await compressed_stream.read(chunk_size) + if not chunk: + break + + # Decompress chunk and yield + yield decompressor.decompress(chunk) + + logger.info(f"Streamed and decompressed {file_name} from {key_name}") + + except ClientError as e: + logger.error(f"S3 error streaming file {key_name}: {e}") + yield f"Error: {str(e)}".encode('utf-8') + except Exception as e: + logger.error(f"Error streaming and decompressing {key_name}: {e}") + yield f"Error: {str(e)}".encode('utf-8') + + return content_stream() From 368e4d99f264dfa070587149a823da8b85aaf974 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:07:39 +0300 Subject: [PATCH 21/30] refactor(utils): improve task status response and error handling - Fix logic flow in task status handling and optimize response creation --- utils.py | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/utils.py b/utils.py index 4f36fec..eeea8ba 100644 --- a/utils.py +++ b/utils.py @@ -234,7 +234,7 @@ def convert_celery_status(celery_status: CeleryTaskStatus) -> TaskStatus: return status_mapping.get(celery_status, TaskStatus.READY) # Default to READY if status is unknown -def create_task_status_response(celery_task:AsyncResult, task: Dict[str, str], task_id: str, base_url: str) -> dict: +def create_task_status_response(celery_task: AsyncResult, task: Dict[str, str], task_id: str, base_url: str) -> dict: """Create response for task status check.""" response = { "task_id": task_id, @@ -248,16 +248,20 @@ def create_task_status_response(celery_task:AsyncResult, task: Dict[str, str], t } if task["status"] == TaskStatus.COMPLETED or celery_task.ready(): - # Safely parse result if it exists and isn't empty - if celery_task.successful() and celery_task.result: + # Handle successful tasks + if celery_task.successful(): + # Always prioritize Celery result, even if it's None or empty response["result"] = celery_task.result - elif task.get("result") and task["result"].strip(): + # Only fall back to Redis result if Celery result is not available + elif not hasattr(celery_task, 'result') and task.get("result"): try: + # Try to parse Redis task result as JSON response["result"] = json.loads(task["result"]) except json.JSONDecodeError: - # If we can't parse the result as JSON, use it as a string + # If parsing fails, use it as a string response["result"] = task["result"] else: + # Set explicit None if no result is available response["result"] = None elif task["status"] == TaskStatus.FAILED or celery_task.failed(): response["error"] = task.get("error", "Unknown error") From b3db86da799a4338c6bda2a23fc8d5ba31b3865e Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:08:03 +0300 Subject: [PATCH 22/30] fix(redisCache): enhance connection handling and URL parsing - Improve Redis connection reliability with better error handling, URL parsing, and client type detection --- redisCache.py | 48 ++++++++++++++++++++++++++++++++---------------- 1 file changed, 32 insertions(+), 16 deletions(-) diff --git a/redisCache.py b/redisCache.py index e6d56f6..5d9db28 100644 --- a/redisCache.py +++ b/redisCache.py @@ -7,6 +7,7 @@ from upstash_redis.asyncio import Redis from redis.asyncio import Redis as PureRedis import asyncio +from urllib.parse import urlparse REDIS_CHANNEL = "stream_channel" # Default channel for streaming data @@ -24,10 +25,20 @@ REDIS_USERNAME = os.environ.get("UPSTASH_REDIS_USER") REDIS_PASSWORD = os.environ.get("UPSTASH_REDIS_PASS") +# not all([redis_url, redis_token, REDIS_PORT, REDIS_USERNAME, REDIS_PASSWORD]): if not redis_url or not redis_token or not REDIS_PORT or not REDIS_USERNAME or not REDIS_PASSWORD: - raise ValueError("UPSTASH_REDIS_REST_URL and UPSTASH_REDIS_REST_TOKEN environment variables must be set") - -REDIS_URL = redis_url.replace("https://", "") + missing = [ + name for name, val in [ + ("UPSTASH_REDIS_REST_URL", redis_url), + ("UPSTASH_REDIS_REST_TOKEN", redis_token), + ("UPSTASH_REDIS_PORT", REDIS_PORT), + ("UPSTASH_REDIS_USER", REDIS_USERNAME), + ("UPSTASH_REDIS_PASS", REDIS_PASSWORD), + ] if not val + ] + raise ValueError(f"Missing required Redis env vars: {', '.join(missing)}") + +REDIS_URL = urlparse(redis_url).hostname or redis_url.replace("https://", "").replace("http://", "") # ────────────────── redis client ────────────────── # Initialize Redis client with Upstash credentials @@ -50,8 +61,8 @@ async def test_connection(redis: Redis | PureRedis): retries = 0 - max_retries = 100 - retry_delay = 12 # seconds + max_retries = int(os.getenv("REDIS_MAX_RETRIES", "10")) + retry_delay = float(os.getenv("REDIS_RETRY_DELAY_SEC", "3.0")) # seconds while retries < max_retries: try: @@ -89,17 +100,21 @@ async def test_connection(redis: Redis | PureRedis): # ─────────────────── redis execute ────────────────── -async def redis_execute(redis: Redis, command: List, *args): +async def redis_execute(redis: Redis | PureRedis, command: List, *args): """Execute a Redis command and handle errors.""" if not redis: print("\033[91mERROR-DB:\033[0m Redis client is not initialized.") return None - if not command: - print("\033[91mERROR-DB:\033[0m Command is empty.") - return None try: - result = await redis.execute(command, *args) + if not command: + print("\033[91mERROR-DB:\033[0m Command is empty.") + raise ValueError("command must be a non-empty list") + + if isinstance(redis, Redis) and hasattr(redis, "execute"): + result = await redis.execute(command, *args) + elif isinstance(redis, PureRedis) and hasattr(redis, "execute_command"): + result = await redis.execute_command(*command, *args) return result except Exception as e: print(f"\033[91mERROR-DB:\033[0m Redis command '{command}' failed: {e}") @@ -159,10 +174,11 @@ async def redis_xadd(pipe, channel: str, message: dict, maxlen: Optional[int] = message_id = await pipe.xadd(channel, message, **kwargs) else: # Fallback: Upstash `execute` or redis-py `execute_command` - if hasattr(pipe, "execute"): + # Flatten into varargs for the underlying client + if isinstance(redis, Redis) and hasattr(pipe, "execute"): message_id = await pipe.execute("XADD", *pieces) - else: - message_id = await pipe.execute_command("XADD", *pieces) + elif isinstance(redis, PureRedis) and hasattr(pipe, "execute_command"): + message_id = await pipe.execute_command("XADD", *pieces) return message_id except Exception as e: print(f"\033[91mERROR-DB:\033[0m Failed to add message to stream '{channel}': {e}") @@ -198,9 +214,9 @@ async def redis_xread(redis: Redis | PureRedis, streams: dict, count: Optional[i for channel, last_id in streams.items(): command.append(last_id) # Flatten into varargs for the underlying client - if hasattr(redis, "execute"): - result = await redis.execute(command[0], *command[1:]) - else: + if isinstance(redis, Redis) and hasattr(redis, "execute"): + result = await redis.execute(command) + elif isinstance(redis, PureRedis) and hasattr(redis, "execute_command"): result = await redis.execute_command(*command) return result except Exception as e: From 0fd2460fd10ca1bfa74423ea9785e3ecbb1aaf86 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:08:23 +0300 Subject: [PATCH 23/30] refactor(api): enhance error handling and type safety in streaming functions - Add proper type casting, improve error handling, and fix SSE formatting in streaming responses --- api.py | 74 ++++++++++++++++++++++++++++++++++++---------------------- 1 file changed, 46 insertions(+), 28 deletions(-) diff --git a/api.py b/api.py index 63e7dd9..d3ac72e 100644 --- a/api.py +++ b/api.py @@ -15,7 +15,7 @@ import json import logging import os -from typing import Any, AsyncGenerator, Dict, List, Optional, Tuple +from typing import Any, AsyncGenerator, Dict, List, Optional, Tuple, cast from urllib.parse import unquote from celery.result import AsyncResult # Import AsyncResult here from celery_app import celery_app # Import celery_app here @@ -328,30 +328,36 @@ async def handle_stream_crawl_request( ) -> Tuple[AsyncWebCrawler, AsyncGenerator]: """Handle streaming crawl requests.""" try: - browser_config = BrowserConfig.load(browser_config) + browser_conf = BrowserConfig.load(browser_config) # browser_config.verbose = True # Set to False or remove for production stress testing - browser_config.verbose = False - crawler_config = CrawlerRunConfig.load(crawler_config) - crawler_config.scraping_strategy = LXMLWebScrapingStrategy() - crawler_config.stream = True + browser_conf.verbose = False + crawler_conf = CrawlerRunConfig.load(crawler_config) + crawler_conf.scraping_strategy = LXMLWebScrapingStrategy() + crawler_conf.stream = True + + crawler_cfg = (config.get("crawler") or {}) + crl_cfg = (config.get("rate_limiter") or {}) dispatcher = MemoryAdaptiveDispatcher( - memory_threshold_percent=config["crawler"]["memory_threshold_percent"], - rate_limiter=RateLimiter( - base_delay=tuple(config["crawler"]["rate_limiter"]["base_delay"]) - ) + memory_threshold_percent=crawler_cfg.get("memory_threshold_percent", 80), + rate_limiter=RateLimiter( + base_delay=tuple(crl_cfg.get("base_delay", (0.2, 1.0))) + ) if crl_cfg.get("enabled", False) else None ) from crawler_pool import get_crawler - bcrawler:tuple[AsyncWebCrawler, str] = await get_crawler(browser_config) + bcrawler:tuple[AsyncWebCrawler, str] = await get_crawler(browser_conf) crawler, _ = bcrawler # crawler = AsyncWebCrawler(config=browser_config) # await crawler.start() - results_gen = await crawler.arun_many( + results_gen: AsyncGenerator = cast( + AsyncGenerator, + await crawler.arun_many( urls=urls, - config=crawler_config, + config=crawler_conf, dispatcher=dispatcher + ) ) return crawler, results_gen @@ -487,10 +493,14 @@ async def cancel_a_job( response = {"status": "ok", "message": f"Task {temp_task_id} canceled successfully."} task_id = await redis.hget(key=f"temp_task_id:{temp_task_id}", field='celery_task_id') + if isinstance(task_id, bytes): + task_id = task_id.decode('utf-8') - task_info = await redis.hgetall(f"task:{task_id}") + task_raw = await redis.hgetall(f"task:{task_id}") - operation_id = task_info.get("operation_id") if task_info else None + task_info = decode_redis_hash(task_raw) if task_raw else {} + + operation_id = task_info.get("operation_id") if not task_info: raise HTTPException(status_code=404, detail="Task not found") @@ -518,7 +528,7 @@ async def cancel_a_job( terminate=True, signal=signal.SIGKILL if force else signal.SIGTERM ) - + deadline = time.monotonic() + 15 # seconds while True: # Check the task status to see if it's finished if celery_task.ready(): @@ -582,6 +592,10 @@ async def cancel_a_job( ) break + + if time.monotonic() > deadline: + logger.warning("Timeout waiting for task %s to cancel", task_id) + break await asyncio.sleep(0.3) # Wait before checking again @@ -615,16 +629,20 @@ async def handle_task_status( # query celery task state by task id task = decode_redis_hash(task) - + temp_task_id = task.get("temp_task_id") response = create_task_status_response(celery_task, task, task_id, base_url) # remove task from redis keep metadata in firebase - if task["status"] in [TaskStatus.COMPLETED, TaskStatus.FAILED, TaskStatus.CANCELED]: + status_str = response["status"].value if isinstance(response["status"], TaskStatus) else response["status"] + if status_str in {TaskStatus.COMPLETED.value, TaskStatus.FAILED.value, TaskStatus.CANCELED.value}: if not keep and should_cleanup_task(task["created_at"]): - await redis.delete(f"task:{task_id}") - await redis.delete(f"{REDIS_CHANNEL}:{task_id}") - await redis.delete(f"celery-task-meta-{task_id}") - await redis.delete(f"temp_task_id:{task['temp_task_id']}") + pipe = redis.multi() + pipe.delete(f"task:{task_id}") + pipe.delete(f"{REDIS_CHANNEL}:{task_id}") + pipe.delete(f"celery-task-meta-{task_id}") + if temp_task_id: + pipe.delete(f"temp_task_id:{temp_task_id}") + await pipe.exec() return JSONResponse(response) @@ -652,7 +670,7 @@ async def stream_task_status(): # Generate status response response = create_task_status_response(celery_task, task, task_id, base_url) - data = f"data: {json.dumps(response)}\n" # Note the double newline for SSE format + data = f"data: {json.dumps(response)}\n\n" # Note the double newline for SSE format yield data.encode('utf-8') # Ensure we're yielding bytes @@ -665,19 +683,19 @@ async def stream_task_status(): # Handle any serialization errors error_msg = f"Error generating status response: {str(e)}" logger.error(error_msg) - yield f"data: {json.dumps({'error': 'An internal error occurred while generating the status response.'})}\n".encode('utf-8') + yield f"data: {json.dumps({'error': error_msg})}\n".encode('utf-8') # Wait before checking again await asyncio.sleep(1) # Send the [DONE] marker to end the stream - yield b"data: [DONE]\n" + yield b"data: [DONE]\n\n" except Exception as e: # TODO: Handle exceptions in the streaming loop IN the frontend logger.error(f"Fatal error in status stream: {str(e)}", exc_info=True) - yield f"event: error\ndata: {json.dumps({'error': 'A fatal error occurred while streaming the task status.', 'fatal': True})}\n".encode('utf-8') - yield b"data: [DONE]\n" + yield f"event: error\ndata: {json.dumps({'error': 'A fatal error occurred while streaming the task status.', 'fatal': True})}\n\n".encode('utf-8') + yield b"data: [DONE]\n\n" return StreamingResponse( stream_task_status(), @@ -691,7 +709,7 @@ async def stream_task_status(): ) except Exception as e: # Return a proper error response instead of letting the exception bubble up - logger.error(f"Error setting up status stream for task {task_id}: {str(e)}", exc_info=True) + logger.exception("Error setting up status stream for task %s: %s", task_id, e) return JSONResponse( status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, content={ From 4d7cb33ccb9caa14a3bca25ccecc4adab64be094 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:08:43 +0300 Subject: [PATCH 24/30] fix(job): improve SSE message formatting and error handling - Fix SSE message format with proper double newlines and improve error handling in stream results --- job.py | 23 ++++++++++++----------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/job.py b/job.py index f5c131c..9627e1d 100644 --- a/job.py +++ b/job.py @@ -428,7 +428,7 @@ async def event_stream(channel: str): # yield b"data: {\"message\":\"completed\",\"type\":\"auto_complete\"}\n\n" # yield b"data: [DONE]\n\n" break - elif retries > max_retries and not completed_yielded and celery_task.state == "PENDING" or celery_task.state == "STARTED": + elif retries > max_retries and not completed_yielded and (celery_task.state in {"PENDING", "STARTED"}): retries = 0 @@ -462,9 +462,9 @@ async def event_stream(channel: str): if not completed_yielded: logger.info(f"Task {task_id}: Yielding completion message") # ADD 'data: ' PREFIX HERE - yield f"data: {json.dumps(msg_data_dict, ensure_ascii=False)}\n".encode('utf-8') + yield f"data: {json.dumps(msg_data_dict, ensure_ascii=False)}\n\n".encode('utf-8') completed_yielded = True - yield b"data: [DONE]\n" + yield b"data: [DONE]\n\n" return continue @@ -484,16 +484,17 @@ async def event_stream(channel: str): seen_messages.add(unique_id) # ADD 'data: ' PREFIX HERE - yield f"data: {json.dumps(msg_data_dict, ensure_ascii=False)}\n".encode('utf-8') + yield f"data: {json.dumps(msg_data_dict, ensure_ascii=False)}\n\n".encode('utf-8') else: # No messages - increment retry counter retries += 1 logger.warning(f"No messages returned or malformed response. Retry count: {retries}") except Exception as e: - logger.error(f"Error reading from Redis stream: {e}") + logger.exception("Error reading from Redis stream") # Send error to client as an SSE event - yield f"event: error\ndata: {json.dumps({'error': str(e)})}\n".encode('utf-8') + # Send generic error to client as an SSE event + yield f"event: error\ndata: {json.dumps({'error': 'stream read error'})}\n\n".encode('utf-8') # Optionally end the stream after a serious error retries += 1 @@ -501,17 +502,17 @@ async def event_stream(channel: str): await asyncio.sleep(poll_interval) except asyncio.CancelledError: - logger.info(f"Task {task_id}: Stream cancelled by client") - yield b"data: {\"message\":\"stream_cancelled\"}\n" + logger.info("Task %s: Stream cancelled by client", task_id) + yield b"event: canceled\ndata: {\"message\":\"stream_cancelled\"}\n\n" except Exception as e: - logger.error(f"Fatal error in event stream: {e}", exc_info=True) - yield f"event: error\ndata: {json.dumps({'error': str(e), 'fatal': True})}\n".encode('utf-8') + logger.exception("Fatal error in event stream") + yield f"event: error\ndata: {json.dumps({'error': 'fatal stream error', 'fatal': True})}\n\n".encode('utf-8') finally: seen_messages.clear() # Always send DONE if we exit the loop without returning if not completed_yielded: - yield b"data: [DONE]\n" + yield b"data: [DONE]\n\n" logger.info(f"Task {task_id}: Stream closed") return StreamingResponse( From fcfe6756eff0992e9e7dbd3f8793133d30afd9dd Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:10:51 +0300 Subject: [PATCH 25/30] chore(Dockerfile): update entry point and simplify CMD - Simplify CMD by moving noVNC handling to the entrypoint script for better process management --- Dockerfile | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/Dockerfile b/Dockerfile index 7a6ce9f..375dbcc 100644 --- a/Dockerfile +++ b/Dockerfile @@ -119,7 +119,7 @@ COPY --chown=appuser:appuser . . COPY --chown=appuser:appuser config.yml . # Set display environment variable -ENV DISPLAY=:99 +# ENV DISPLAY=:99 # Run diagnostics RUN crawl4ai-doctor @@ -142,6 +142,4 @@ USER appuser ENTRYPOINT ["docker-entrypoint.sh"] # Start application -CMD ["sh", "-c", "\ - /opt/noVNC/utils/websockify/run --web /opt/noVNC 0.0.0.0:6080 0.0.0.0:5900 & \ - uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] +CMD ["sh", "-c", "uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] From dca97c90a956173fda73af63e6f6e3c0428990f8 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:10:58 +0300 Subject: [PATCH 26/30] chore(fly.toml): update service configuration and add VNC security - Add VNC password security, optimize service definitions, and update process commands --- fly.toml | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/fly.toml b/fly.toml index 72e8b07..3a6fc3f 100644 --- a/fly.toml +++ b/fly.toml @@ -7,14 +7,20 @@ app = 'crawlagent' primary_region = 'fra' [build] + # Add build configuration for better build process + dockerfile = "Dockerfile" + # build-target = "final" [env] PYTHON_ENV = 'production' UPSTASH_REDIS_PORT = '30766' UPSTASH_REDIS_REST_URL = 'https://gusc1-saved-terrapin-30766.upstash.io' + REDIS_MAX_RETRIES = "100" + REDIS_RETRY_DELAY_SEC = "3" AWS_ENDPOINT_URL_S3 = 'https://fly.storage.tigris.dev' AWS_ENDPOINT_URL_IAM = 'https://fly.iam.storage.tigris.dev' AWS_REGION = 'auto' + VNC_PASSWORD = "${VNC_PASSWORD}" # Set this as a secret [[services]] protocol = 'tcp' @@ -50,15 +56,14 @@ primary_region = 'fra' [[services]] internal_port = 6080 protocol = "tcp" - processes = ["app","worker"] + processes = ["app"] # , "worker" Critical: Remove public noVNC (6080) or enforce VNC auth + restrict binding [[services.ports]] port = 6080 handlers = ["tls", "http"] # Ensure "tls" is included for HTTPS [processes] - app = "sh -c '/opt/noVNC/utils/websockify/run --web /opt/noVNC 0.0.0.0:6080 0.0.0.0:5900 & uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets'" - worker = "celery -A celery_app.celery_app worker --loglevel=info" - + app = "exec uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets" + worker = "exec celery -A celery_app.celery_app worker --loglevel=info" # app = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets'" # worker = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && celery -A celery_app.celery_app worker --loglevel=info --concurrency=2'" From 2c22595b3ce03c03196c0d9e5e2e9b36c64ba2f6 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Tue, 23 Sep 2025 23:11:15 +0300 Subject: [PATCH 27/30] refactor(celery_app): update Redis connection handling - Improve URL parsing and environment variable management for Celery Redis connection --- celery_app.py | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/celery_app.py b/celery_app.py index d94bbdd..d181a1e 100644 --- a/celery_app.py +++ b/celery_app.py @@ -1,12 +1,11 @@ -import asyncio import logging import sys -from kombu import Queue +# from kombu import Queue import os from celery import Celery from dotenv import load_dotenv import signal - +from urllib.parse import urlparse # if sys.platform == "win32": # asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()) @@ -23,16 +22,16 @@ # For Upstash Redis, the URL is typically what's needed. # We'll use the URL from redisCache.py's environment variables.UPSTASH_REDIS_REST_PASSWORD # rediss://default:4c0962711bf64ff8b7797d38dc0e69e5@gusc1-saved-terrapin-30766.upstash.io:30766/0?ssl_cert_reqs=CERT_REQUIRED -REDIS_URL = os.environ.get("UPSTASH_REDIS_REST_URL") +redis_url = os.environ.get("UPSTASH_REDIS_REST_URL") REDIS_PORT = os.environ.get("UPSTASH_REDIS_PORT") REDIS_USERNAME = os.environ.get("UPSTASH_REDIS_USER") REDIS_PASSWORD = os.environ.get("UPSTASH_REDIS_PASS") -if not REDIS_URL or not REDIS_PORT or not REDIS_PASSWORD or not REDIS_USERNAME: - raise ValueError("UPSTASH_REDIS_REST_URL, UPSTASH_REDIS_PORT, UPSTASH_REDIS_USER, UPSTASH_REDIS_REST_PASSWORD environment variables must be set for Celery configuration.") +if not redis_url or not REDIS_PORT or not REDIS_PASSWORD or not REDIS_USERNAME: + raise ValueError("UPSTASH_REDIS_REST_URL, UPSTASH_REDIS_PORT, UPSTASH_REDIS_USER, UPSTASH_REDIS_PASS environment variables must be set for Celery configuration.") -REDIS_URL = REDIS_URL.replace("https://", "") +REDIS_URL = urlparse(redis_url).hostname or redis_url.replace("https://", "").replace("http://", "") REDIS_URI = f"rediss://{REDIS_USERNAME}:{REDIS_PASSWORD}@{REDIS_URL}:{REDIS_PORT}/0?ssl_cert_reqs=CERT_REQUIRED" From cf444ea7e70aed34d749fbb0a796b9ccdc58cd5b Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Wed, 24 Sep 2025 01:46:34 +0300 Subject: [PATCH 28/30] chore(build): Enhance build process with UV and refine Dockerfile This commit significantly refactors the Docker build process by integrating `uv` for Python dependency management. These changes aim to improve build performance, consistency, and maintainability of the Docker images. Key changes include: * **Introduced `uv`:** Replaced `pip` with `uv` for installing Python packages, leveraging its speed and reliability for dependency resolution and installation. This involves adding a dedicated `uv` build stage and updating installation commands in both builder and final stages. * **Streamlined dependency installation:** Removed `PIP_*` environment variables and simplified the wheel building process. System-level build dependencies are now explicitly listed in the final stage for robustness. * **Refined Dockerfile structure:** Updated environment variables, adjusted the `docker-entrypoint.sh` handling for better clarity, and switched to direct `uvicorn` execution in the `CMD`. * **Updated `fly.toml`:** Removed the `AWS_ENDPOINT_URL_S3` configuration, which is no longer needed. * Make sure the file has LF line endings (\n), not CRLF (\r\n), otherwise No such file or directory can also appear. Ensure #!/bin/sh at the top of the script is correct and sh exists in the base image. --- Dockerfile | 88 +++++---- fly.toml | 6 +- pyproject.toml | 83 +++++++++ uv.lock | 487 +++++++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 627 insertions(+), 37 deletions(-) create mode 100644 pyproject.toml create mode 100644 uv.lock diff --git a/Dockerfile b/Dockerfile index 375dbcc..35d754f 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,46 +1,63 @@ -# Build stage +# Build stage with UV +FROM ghcr.io/astral-sh/uv:0.8.22 AS uv + +# Builder stage FROM python:3.12-slim AS builder # Set build-time environment variables ENV PYTHONFAULTHANDLER=1 \ PYTHONHASHSEED=random \ PYTHONUNBUFFERED=1 \ - PIP_NO_CACHE_DIR=1 \ PYTHONDONTWRITEBYTECODE=1 \ - PIP_DISABLE_PIP_VERSION_CHECK=1 \ - PIP_DEFAULT_TIMEOUT=100 \ - DEBIAN_FRONTEND=noninteractive + DEBIAN_FRONTEND=noninteractive \ + UV_COMPILE_BYTECODE=1 \ + UV_NO_INSTALLER_METADATA=1 \ + UV_LINK_MODE=copy -# Set build arguments ARG APP_HOME=/app - -# Create app directory WORKDIR ${APP_HOME} # Install build dependencies COPY requirements.txt . RUN apt-get update && apt-get install -y --no-install-recommends \ - build-essential gcc g++ python3-dev pkg-config libjpeg-dev cmake \ - && rm -rf /var/lib/apt/lists/* \ - && pip wheel --no-cache-dir --no-deps --wheel-dir /app/wheels -r requirements.txt + build-essential \ + curl \ + gcc \ + g++ \ + python3-dev \ + pkg-config \ + libjpeg-dev \ + cmake \ + && rm -rf /var/lib/apt/lists/* + +# Use UV to build wheels +RUN --mount=from=uv,source=/uv,target=/bin/uv \ + --mount=type=cache,target=/root/.cache/uv \ + uv pip install --system -r requirements.txt + # Final stage FROM python:3.12-slim # Metadata LABEL maintainer="Prokopis Antoniadis" \ description="🔥🕷️ Crawl4AI: LLM Web Crawler & scraper" \ - version="1.0" + version="0.1.0" # Set environment variables ENV PYTHONUNBUFFERED=1 \ PYTHONDONTWRITEBYTECODE=1 \ PLAYWRIGHT_BROWSERS_PATH=/ms-playwright \ PYTHON_ENV=production \ + UV_LINK_MODE=copy \ DISPLAY=:99 # Set build arguments ARG APP_HOME=/app +# Install UV in final stage +COPY --from=uv /uv /usr/local/bin/uv +RUN chmod +x /usr/local/bin/uv + WORKDIR ${APP_HOME} # Install system dependencies in a single layer @@ -49,18 +66,17 @@ RUN apt-get update && apt-get install -y --no-install-recommends \ ca-certificates \ lsof \ # Add build dependencies for madoka - build-essential \ + build-essential \ curl \ - wget \ - gnupg \ - cmake \ - gcc \ - g++ \ - pkg-config \ + gcc \ + g++ \ python3-dev \ + pkg-config \ libjpeg-dev \ + cmake \ + wget \ + gnupg \ supervisor \ - libnss3 \ # Playwright system dependencies libglib2.0-0 \ libnss3 \ @@ -102,17 +118,19 @@ RUN groupadd -r appuser && \ mkdir -p /home/appuser/.cache /ms-playwright && \ chown -R appuser:appuser /home/appuser /ms-playwright ${APP_HOME} -# Copy wheels from builder -COPY --from=builder /app/wheels /wheels -COPY --from=builder /app/requirements.txt . +# Install Python dependencies using UV +# COPY --from=builder /app/wheels /wheels + +# Copy dependencies and install +COPY --from=builder ${APP_HOME}/requirements.txt . +RUN --mount=type=cache,target=/root/.cache/uv \ + uv pip install --system -r requirements.txt && \ + uv pip install --system playwright websockets && \ + playwright install --with-deps chromium -# Install dependencies -RUN pip install --no-cache-dir /wheels/* \ - && rm -rf /wheels \ - && pip install --no-cache-dir playwright \ - && playwright install --with-deps chromium \ - && python -c "import crawl4ai; print('✅ crawl4ai is ready to rock!')" \ - && python -c "from playwright.sync_api import sync_playwright; print('✅ Playwright is feeling dramatic!')" +# Verify installations +RUN python -c "import crawl4ai; print('✅ crawl4ai is ready to rock!')" && \ + python -c "from playwright.sync_api import sync_playwright; print('✅ Playwright is feeling dramatic!')" # Copy application code COPY --chown=appuser:appuser . . @@ -132,14 +150,16 @@ EXPOSE 8000 9222 6080 # CMD curl -f http://localhost:8000/health || exit 1 # Copy and set permissions for the entrypoint script -COPY --chown=root:root docker-entrypoint.sh /usr/local/bin/docker-entrypoint.sh -RUN chmod +x /usr/local/bin/docker-entrypoint.sh +COPY docker-entrypoint.sh /usr/local/bin/ +RUN chmod +x /usr/local/bin/docker-entrypoint.sh && \ + chown root:root /usr/local/bin/docker-entrypoint.sh && \ + ls -la /usr/local/bin/docker-entrypoint.sh # Verify permissions # Switch to non-root user USER appuser # Set the entrypoint -ENTRYPOINT ["docker-entrypoint.sh"] +ENTRYPOINT ["/usr/local/bin/docker-entrypoint.sh"] # Start application -CMD ["sh", "-c", "uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets"] +CMD ["uvicorn", "server:app", "--host", "0.0.0.0", "--port", "8000", "--ws", "websockets"] diff --git a/fly.toml b/fly.toml index 3a6fc3f..f6d2618 100644 --- a/fly.toml +++ b/fly.toml @@ -20,7 +20,6 @@ primary_region = 'fra' AWS_ENDPOINT_URL_S3 = 'https://fly.storage.tigris.dev' AWS_ENDPOINT_URL_IAM = 'https://fly.iam.storage.tigris.dev' AWS_REGION = 'auto' - VNC_PASSWORD = "${VNC_PASSWORD}" # Set this as a secret [[services]] protocol = 'tcp' @@ -62,8 +61,9 @@ primary_region = 'fra' handlers = ["tls", "http"] # Ensure "tls" is included for HTTPS [processes] - app = "exec uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets" - worker = "exec celery -A celery_app.celery_app worker --loglevel=info" + app = "uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets" + worker = "celery -A celery_app.celery_app worker --loglevel=info" + # app = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && uvicorn server:app --host 0.0.0.0 --port 8000 --ws websockets'" # worker = "sh -c 'Xvfb :99 -screen 0 1280x720x24 & export DISPLAY=:99 && celery -A celery_app.celery_app worker --loglevel=info --concurrency=2'" diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..b32a16b --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,83 @@ +[project] +name = "crawlagent" +version = "0.1.0" +description = "🔥🕷️ Crawl4AI: LLM Web Crawler & scraper" +authors = [ + { name = "Prokopis Antoniadis", email = "prokopis123@gmail.com" } +] +dependencies = [ + # Web Framework and Server + "fastapi>=0.100.0", + "uvicorn[standard]>=0.22.0", + "gunicorn>=21.2.0", + "websockets>=11.0.3", + "uvloop>=0.17.0", + + # Crawling and Scraping + "crawl4ai>=0.1.0", + "beautifulsoup4~=4.12", + "tf-playwright-stealth>=1.1.0", + "courlan>=0.9.0", + + # Database and Caching + "firebase-admin>=6.2.0", + "google-cloud-firestore>=2.11.0", + "upstash-redis~=1.4.0", + "upstash_ratelimit>=0.0.7", + "redis[hiredis]>=4.6.0", + "celery>=5.3.1", + + # AWS and Storage + "aioboto3>=11.2.0", + "zstandard>=0.21.0", + + # Utils and Monitoring + "python-dotenv>=1.0.0", + "pydantic>=2.10", + "psutil>=6.1.1", + "aiomultiprocess>=0.9.0", + "apscheduler>=3.10.0", + "prometheus_client>=0.17.0", + "prometheus-fastapi-instrumentator>=6.1.0", +] + +[project.optional-dependencies] +dev = [ + "pytest>=7.4.0", + "pytest-asyncio>=0.21.1", + "black>=23.7.0", + "isort>=5.12.0", + "mypy>=1.4.1", + "ruff>=0.0.280", +] + +[build-system] +requires = ["hatchling"] +build-backend = "hatchling.build" + +[tool.uv] +target-version = ["py312"] +resolve-log = true +compile-bytecode = true +no-installer-metadata = true +link-mode = "copy" + +[tool.black] +line-length = 88 +target-version = ["py312"] + +[tool.isort] +profile = "black" +multi_line_output = 3 + +[tool.mypy] +python_version = "3.12" +strict = true +warn_return_any = true +warn_unused_configs = true +disallow_untyped_defs = true + +[tool.ruff] +line-length = 88 +target-version = "py312" +select = ["E", "F", "B", "I"] \ No newline at end of file diff --git a/uv.lock b/uv.lock new file mode 100644 index 0000000..2830f6b --- /dev/null +++ b/uv.lock @@ -0,0 +1,487 @@ +# This file was autogenerated by uv via the following command: +# uv pip compile pyproject.toml -o uv.lock +aioboto3==15.1.0 + # via crawlagent (pyproject.toml) +aiobotocore==2.24.0 + # via aioboto3 +aiofiles==24.1.0 + # via + # aioboto3 + # crawl4ai +aiohappyeyeballs==2.6.1 + # via aiohttp +aiohttp==3.12.15 + # via + # aiobotocore + # crawl4ai + # litellm +aioitertools==0.12.0 + # via aiobotocore +aiomultiprocess==0.9.1 + # via crawlagent (pyproject.toml) +aiosignal==1.4.0 + # via aiohttp +aiosqlite==0.21.0 + # via crawl4ai +alphashape==1.3.1 + # via crawl4ai +amqp==5.3.1 + # via kombu +annotated-types==0.7.0 + # via pydantic +anyio==4.11.0 + # via + # crawl4ai + # httpx + # openai + # starlette + # watchfiles +apscheduler==3.11.0 + # via crawlagent (pyproject.toml) +async-timeout==5.0.1 + # via + # aiohttp + # redis +attrs==25.3.0 + # via + # aiohttp + # jsonschema + # referencing +babel==2.17.0 + # via courlan +beautifulsoup4==4.13.5 + # via + # crawlagent (pyproject.toml) + # crawl4ai +billiard==4.2.2 + # via celery +boto3==1.39.11 + # via aiobotocore +botocore==1.39.11 + # via + # aiobotocore + # boto3 + # s3transfer +brotli==1.1.0 + # via crawl4ai +cachecontrol==0.14.3 + # via firebase-admin +cachetools==5.5.2 + # via google-auth +celery==5.5.3 + # via crawlagent (pyproject.toml) +certifi==2025.8.3 + # via + # httpcore + # httpx + # requests +cffi==2.0.0 + # via cryptography +chardet==5.2.0 + # via crawl4ai +charset-normalizer==3.4.3 + # via requests +click==8.3.0 + # via + # alphashape + # celery + # click-didyoumean + # click-log + # click-plugins + # click-repl + # crawl4ai + # litellm + # nltk + # uvicorn +click-didyoumean==0.3.1 + # via celery +click-log==0.4.0 + # via alphashape +click-plugins==1.1.1.2 + # via celery +click-repl==0.3.0 + # via celery +courlan==1.3.2 + # via crawlagent (pyproject.toml) +crawl4ai==0.7.4 + # via crawlagent (pyproject.toml) +cryptography==46.0.1 + # via + # pyjwt + # pyopenssl +distro==1.9.0 + # via openai +exceptiongroup==1.3.0 + # via anyio +fake-http-header==0.3.5 + # via tf-playwright-stealth +fake-useragent==2.2.0 + # via crawl4ai +fastapi==0.117.1 + # via crawlagent (pyproject.toml) +fastuuid==0.12.0 + # via litellm +filelock==3.19.1 + # via huggingface-hub +firebase-admin==7.1.0 + # via crawlagent (pyproject.toml) +frozenlist==1.7.0 + # via + # aiohttp + # aiosignal +fsspec==2025.9.0 + # via huggingface-hub +google-api-core==2.25.1 + # via + # firebase-admin + # google-cloud-core + # google-cloud-firestore + # google-cloud-storage +google-auth==2.40.3 + # via + # google-api-core + # google-cloud-core + # google-cloud-firestore + # google-cloud-storage +google-cloud-core==2.4.3 + # via + # google-cloud-firestore + # google-cloud-storage +google-cloud-firestore==2.21.0 + # via + # crawlagent (pyproject.toml) + # firebase-admin +google-cloud-storage==3.4.0 + # via firebase-admin +google-crc32c==1.7.1 + # via + # google-cloud-storage + # google-resumable-media +google-resumable-media==2.7.2 + # via google-cloud-storage +googleapis-common-protos==1.70.0 + # via + # google-api-core + # grpcio-status +greenlet==3.2.4 + # via + # patchright + # playwright +grpcio==1.75.0 + # via + # google-api-core + # grpcio-status +grpcio-status==1.75.0 + # via google-api-core +gunicorn==23.0.0 + # via crawlagent (pyproject.toml) +h11==0.16.0 + # via + # httpcore + # uvicorn +h2==4.3.0 + # via httpx +hf-xet==1.1.10 + # via huggingface-hub +hiredis==3.2.1 + # via redis +hpack==4.1.0 + # via h2 +httpcore==1.0.9 + # via httpx +httptools==0.6.4 + # via uvicorn +httpx==0.28.1 + # via + # crawl4ai + # firebase-admin + # litellm + # openai + # upstash-redis +huggingface-hub==0.35.1 + # via tokenizers +humanize==4.13.0 + # via crawl4ai +hyperframe==6.1.0 + # via h2 +idna==3.10 + # via + # anyio + # httpx + # requests + # yarl +importlib-metadata==8.7.0 + # via litellm +jinja2==3.1.6 + # via litellm +jiter==0.11.0 + # via openai +jmespath==1.0.1 + # via + # aiobotocore + # boto3 + # botocore +joblib==1.5.2 + # via nltk +jsonschema==4.25.1 + # via litellm +jsonschema-specifications==2025.9.1 + # via jsonschema +kombu==5.5.4 + # via celery +lark==1.3.0 + # via crawl4ai +litellm==1.77.3 + # via crawl4ai +lxml==5.4.0 + # via crawl4ai +madoka==0.7.1 + # via pondpond +markdown-it-py==4.0.0 + # via rich +markupsafe==3.0.2 + # via jinja2 +mdurl==0.1.2 + # via markdown-it-py +msgpack==1.1.1 + # via cachecontrol +multidict==6.6.4 + # via + # aiobotocore + # aiohttp + # yarl +networkx==3.4.2 + # via alphashape +nltk==3.9.1 + # via crawl4ai +numpy==2.2.6 + # via + # alphashape + # crawl4ai + # rank-bm25 + # scipy + # shapely + # trimesh +openai==1.109.0 + # via litellm +packaging==25.0 + # via + # gunicorn + # huggingface-hub + # kombu +patchright==1.55.2 + # via crawl4ai +pillow==11.3.0 + # via crawl4ai +playwright==1.55.0 + # via + # crawl4ai + # tf-playwright-stealth +pondpond==1.4.1 + # via litellm +prometheus-client==0.23.1 + # via + # crawlagent (pyproject.toml) + # prometheus-fastapi-instrumentator +prometheus-fastapi-instrumentator==7.1.0 + # via crawlagent (pyproject.toml) +prompt-toolkit==3.0.52 + # via click-repl +propcache==0.3.2 + # via + # aiohttp + # yarl +proto-plus==1.26.1 + # via + # google-api-core + # google-cloud-firestore +protobuf==6.32.1 + # via + # google-api-core + # google-cloud-firestore + # googleapis-common-protos + # grpcio-status + # proto-plus +psutil==7.1.0 + # via + # crawlagent (pyproject.toml) + # crawl4ai +pyasn1==0.6.1 + # via + # pyasn1-modules + # rsa +pyasn1-modules==0.4.2 + # via google-auth +pycparser==2.23 + # via cffi +pydantic==2.11.9 + # via + # crawlagent (pyproject.toml) + # crawl4ai + # fastapi + # litellm + # openai +pydantic-core==2.33.2 + # via pydantic +pyee==13.0.0 + # via + # patchright + # playwright +pygments==2.19.2 + # via rich +pyjwt==2.10.1 + # via firebase-admin +pyopenssl==25.3.0 + # via crawl4ai +python-dateutil==2.9.0.post0 + # via + # aiobotocore + # botocore + # celery +python-dotenv==1.1.1 + # via + # crawlagent (pyproject.toml) + # crawl4ai + # litellm + # uvicorn +pyyaml==6.0.2 + # via + # crawl4ai + # huggingface-hub + # uvicorn +rank-bm25==0.2.2 + # via crawl4ai +redis==6.4.0 + # via crawlagent (pyproject.toml) +referencing==0.36.2 + # via + # jsonschema + # jsonschema-specifications +regex==2025.9.18 + # via + # nltk + # tiktoken +requests==2.32.5 + # via + # cachecontrol + # crawl4ai + # google-api-core + # google-cloud-storage + # huggingface-hub + # tiktoken +rich==14.1.0 + # via crawl4ai +rpds-py==0.27.1 + # via + # jsonschema + # referencing +rsa==4.9.1 + # via google-auth +rtree==1.4.1 + # via alphashape +s3transfer==0.13.1 + # via boto3 +scipy==1.15.3 + # via alphashape +shapely==2.1.1 + # via + # alphashape + # crawl4ai +six==1.17.0 + # via python-dateutil +sniffio==1.3.1 + # via + # anyio + # openai +snowballstemmer==2.2.0 + # via crawl4ai +soupsieve==2.8 + # via beautifulsoup4 +starlette==0.48.0 + # via + # fastapi + # prometheus-fastapi-instrumentator +tf-playwright-stealth==1.2.0 + # via + # crawlagent (pyproject.toml) + # crawl4ai +tiktoken==0.11.0 + # via litellm +tld==0.13.1 + # via courlan +tokenizers==0.22.1 + # via litellm +tqdm==4.67.1 + # via + # huggingface-hub + # nltk + # openai +trimesh==4.8.2 + # via alphashape +typing-extensions==4.15.0 + # via + # aiosignal + # aiosqlite + # anyio + # beautifulsoup4 + # cryptography + # exceptiongroup + # fastapi + # grpcio + # huggingface-hub + # multidict + # openai + # pydantic + # pydantic-core + # pyee + # pyopenssl + # referencing + # starlette + # typing-inspection + # uvicorn +typing-inspection==0.4.1 + # via pydantic +tzdata==2025.2 + # via kombu +tzlocal==5.3.1 + # via apscheduler +upstash-ratelimit==1.1.0 + # via crawlagent (pyproject.toml) +upstash-redis==1.4.0 + # via + # crawlagent (pyproject.toml) + # upstash-ratelimit +urllib3==2.5.0 + # via + # botocore + # courlan + # requests +uvicorn==0.37.0 + # via crawlagent (pyproject.toml) +uvloop==0.21.0 + # via + # crawlagent (pyproject.toml) + # uvicorn +vine==5.1.0 + # via + # amqp + # celery + # kombu +watchfiles==1.1.0 + # via uvicorn +wcwidth==0.2.14 + # via prompt-toolkit +websockets==15.0.1 + # via + # crawlagent (pyproject.toml) + # uvicorn +wrapt==1.17.3 + # via aiobotocore +xxhash==3.5.0 + # via crawl4ai +yarl==1.20.1 + # via aiohttp +zipp==3.23.0 + # via importlib-metadata +zstandard==0.25.0 + # via crawlagent (pyproject.toml) From c2656b077949d162d800f9730183592b376f4545 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Wed, 24 Sep 2025 02:35:30 +0300 Subject: [PATCH 29/30] chore(gitignore): add regions.md to .gitignore - add regions.md to .gitignore - remove unnecessary websockets installation --- .gitignore | 2 ++ Dockerfile | 2 +- fly.toml | 3 +++ 3 files changed, 6 insertions(+), 1 deletion(-) diff --git a/.gitignore b/.gitignore index bd7b18a..5d30701 100644 --- a/.gitignore +++ b/.gitignore @@ -19,3 +19,5 @@ env/ .vscode/ /*.env + +/*regions.md \ No newline at end of file diff --git a/Dockerfile b/Dockerfile index 35d754f..f71896e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -125,7 +125,7 @@ RUN groupadd -r appuser && \ COPY --from=builder ${APP_HOME}/requirements.txt . RUN --mount=type=cache,target=/root/.cache/uv \ uv pip install --system -r requirements.txt && \ - uv pip install --system playwright websockets && \ + uv pip install --system playwright && \ playwright install --with-deps chromium # Verify installations diff --git a/fly.toml b/fly.toml index f6d2618..8c2f2e8 100644 --- a/fly.toml +++ b/fly.toml @@ -45,6 +45,9 @@ primary_region = 'fra' [[services]] internal_port = 9222 + auto_stop_machines = "stop" + auto_start_machines = true + min_machines_running = 0 protocol = "tcp" processes = ["app", "worker"] From f0c4abf3e6152ac91da1054b9d389167306f7e51 Mon Sep 17 00:00:00 2001 From: prokopis3 Date: Wed, 24 Sep 2025 02:36:56 +0300 Subject: [PATCH 30/30] refactor(storage): replace get_s3_client with S3ClientManager in list_files function --- storage.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/storage.py b/storage.py index 2808ff2..a7cf04d 100644 --- a/storage.py +++ b/storage.py @@ -237,7 +237,7 @@ async def download_file_decompressed(folder_name: str, file_name: str) -> Option async def list_files(folder_name: str) -> List[Dict[str, Any]]: """List all files in a folder.""" try: - async with await get_s3_client() as svc: + async with S3ClientManager() as svc: response = await svc.list_objects_v2( Bucket=TIGRIS_BUCKET_NAME, Prefix=f"{folder_name}/"