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 a42fdef..f71896e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,150 +1,165 @@ -# Use an official Python runtime as a base image -FROM python:3.10-slim AS build +# Build stage with UV +FROM ghcr.io/astral-sh/uv:0.8.22 AS uv -# Set environment variables to production +# 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 \ + UV_COMPILE_BYTECODE=1 \ + UV_NO_INSTALLER_METADATA=1 \ + UV_LINK_MODE=copy + +ARG APP_HOME=/app +WORKDIR ${APP_HOME} + +# Install build dependencies +COPY requirements.txt . +RUN apt-get update && apt-get install -y --no-install-recommends \ + 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="0.1.0" + +# Set environment variables +ENV PYTHONUNBUFFERED=1 \ + PYTHONDONTWRITEBYTECODE=1 \ PLAYWRIGHT_BROWSERS_PATH=/ms-playwright \ - PYTHON_ENV=production + PYTHON_ENV=production \ + UV_LINK_MODE=copy \ + 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 \ + +# 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 +RUN apt-get update && apt-get install -y --no-install-recommends \ + fonts-liberation \ + ca-certificates \ + lsof \ + # Add build dependencies for madoka + build-essential \ + curl \ + gcc \ + g++ \ + python3-dev \ + pkg-config \ + libjpeg-dev \ + cmake \ + wget \ + gnupg \ + supervisor \ + # 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!')" && \ + # 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} + +# 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 && \ + playwright install --with-deps chromium + +# 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 . . +COPY --chown=appuser:appuser config.yml . -# RUN crawl4ai-setup - -# https://playwright.dev/docs/browsers -# Install only the required Playwright browser (Chromium) -RUN playwright install chromium - -# RUN playwright install --with-deps +# 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 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 -# # 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 ["/usr/local/bin/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 ["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"] diff --git a/api.py b/api.py index 6173f9b..d3ac72e 100644 --- a/api.py +++ b/api.py @@ -15,16 +15,17 @@ 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 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 @@ -327,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 @@ -486,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_raw = await redis.hgetall(f"task:{task_id}") - task_info = await redis.hgetall(f"task:{task_id}") + task_info = decode_redis_hash(task_raw) if task_raw else {} - operation_id = task_info.get("operation_id") if task_info else None + operation_id = task_info.get("operation_id") if not task_info: raise HTTPException(status_code=404, detail="Task not found") @@ -504,12 +515,20 @@ 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 + ) + deadline = time.monotonic() + 15 # seconds while True: # Check the task status to see if it's finished if celery_task.ready(): @@ -573,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 @@ -606,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) @@ -624,45 +651,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\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\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\n".encode('utf-8') + yield b"data: [DONE]\n\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.exception("Error setting up status stream for task %s: %s", task_id, e) + return JSONResponse( + status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, + content={ + "status": "error", + "error": "Failed to stream task status due to an internal error." + } + ) async def create_new_task( redis: Redis, diff --git a/celery_app.py b/celery_app.py index af45ce2..d181a1e 100644 --- a/celery_app.py +++ b/celery_app.py @@ -1,45 +1,47 @@ -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()) # 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. # 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" 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 +50,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 +95,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"), 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/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/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/docker-entrypoint.sh b/docker-entrypoint.sh new file mode 100644 index 0000000..3c3b754 --- /dev/null +++ b/docker-entrypoint.sh @@ -0,0 +1,83 @@ +#!/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 +mkdir -p /tmp/.X11-unix +chmod 1777 /tmp/.X11-unix +chown appuser:appuser /tmp/.X11-unix + +# Start Xvfb in the background +Xvfb :99 -screen 0 1280x720x24 -ac & +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 "$@" 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 diff --git a/fly.toml b/fly.toml index fe68f40..8c2f2e8 100644 --- a/fly.toml +++ b/fly.toml @@ -7,10 +7,16 @@ 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' @@ -39,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"] @@ -49,14 +58,16 @@ 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 '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 = "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 +77,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/job.py b/job.py index 5a8d8f2..9627e1d 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,138 @@ 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 in {"PENDING", "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\n".encode('utf-8') + completed_yielded = True + yield b"data: [DONE]\n\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\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.exception("Error reading from Redis stream") + # Send error to client as an SSE event + # 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 + + # Pause before next poll + await asyncio.sleep(poll_interval) - if 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}") + except asyncio.CancelledError: + 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.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\n" + 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/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/redisCache.py b/redisCache.py index cced7bd..5d9db28 100644 --- a/redisCache.py +++ b/redisCache.py @@ -1,17 +1,22 @@ # 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 +from urllib.parse import urlparse 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") @@ -20,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 @@ -41,16 +56,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 = int(os.getenv("REDIS_MAX_RETRIES", "10")) + retry_delay = float(os.getenv("REDIS_RETRY_DELAY_SEC", "3.0")) # 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 ────────────────── @@ -73,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}") @@ -106,4 +137,89 @@ 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)) + # When not specifying an explicit ID, Redis requires "*" + pieces.append("*") + for key, value in message.items(): + pieces.extend([str(key), str(value)]) + + # 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: Upstash `execute` or redis-py `execute_command` + # Flatten into varargs for the underlying client + if isinstance(redis, Redis) and hasattr(pipe, "execute"): + message_id = await pipe.execute("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}") + return 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) + :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) + # Flatten into varargs for the underlying client + 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: + print(f"\033[91mERROR-DB:\033[0m Failed to read from stream(s) '{list(streams.keys())}': {e}") + return None +# ─────────────────── redis pub/sub ────────────────── 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 diff --git a/server.py b/server.py index e6746e5..68dfe9b 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,16 @@ 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__} +@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: + 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"]) diff --git a/storage.py b/storage.py index 63e9145..a7cf04d 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 S3ClientManager() 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() 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..eeea8ba 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 @@ -231,13 +234,13 @@ 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, "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,27 @@ 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"]) + # Handle successful tasks + if celery_task.successful(): + # Always prioritize Celery result, even if it's None or empty + response["result"] = celery_task.result + # 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 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["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 +313,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 +325,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 +335,15 @@ 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() + 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): chunk = batch_json[i:i+chunk_size] @@ -335,27 +353,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 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)