f6f7f1b551
* Fix: Use correct URL variable for raw HTML extraction (#1116) - Prevents full HTML content from being passed as URL to extraction strategies - Added unit tests to verify raw HTML and regular URL processing Fix: Wrong URL variable used for extraction of raw html * Fix #1181: Preserve whitespace in code blocks during HTML scraping The remove_empty_elements_fast() method was removing whitespace-only span elements inside <pre> and <code> tags, causing import statements like "import torch" to become "importtorch". Now skips elements inside code blocks where whitespace is significant. * Refactor Pydantic model configuration to use ConfigDict for arbitrary types * Fix EmbeddingStrategy: Uncomment response handling for the variations and clean up mock data. ref #1621 * Fix: permission issues with .cache/url_seeder and other runtime cache dirs. ref #1638 * fix: ensure BrowserConfig.to_dict serializes proxy_config * feat: make LLM backoff configurable end-to-end - extend LLMConfig with backoff delay/attempt/factor fields and thread them through LLMExtractionStrategy, LLMContentFilter, table extraction, and Docker API handlers - expose the backoff parameter knobs on perform_completion_with_backoff/aperform_completion_with_backoff and document them in the md_v2 guides * reproduced AttributeError from #1642 * pass timeout parameter to docker client request * added missing deep crawling objects to init * generalized query in ContentRelevanceFilter to be a str or list * import modules from enhanceable deserialization * parameterized tests * Fix: capture current page URL to reflect JavaScript navigation and add test for delayed redirects. ref #1268 * refactor: replace PyPDF2 with pypdf across the codebase. ref #1412 * Add browser_context_id and target_id parameters to BrowserConfig Enable Crawl4AI to connect to pre-created CDP browser contexts, which is essential for cloud browser services that pre-create isolated contexts. Changes: - Add browser_context_id and target_id parameters to BrowserConfig - Update from_kwargs() and to_dict() methods - Modify BrowserManager.start() to use existing context when provided - Add _get_page_by_target_id() helper method - Update get_page() to handle pre-existing targets - Add test for browser_context_id functionality This enables cloud services to: 1. Create isolated CDP contexts before Crawl4AI connects 2. Pass context/target IDs to BrowserConfig 3. Have Crawl4AI reuse existing contexts instead of creating new ones * Add cdp_cleanup_on_close flag to prevent memory leaks in cloud/server scenarios * Fix: add cdp_cleanup_on_close to from_kwargs * Fix: find context by target_id for concurrent CDP connections * Fix: use target_id to find correct page in get_page * Fix: use CDP to find context by browserContextId for concurrent sessions * Revert context matching attempts - Playwright cannot see CDP-created contexts * Add create_isolated_context flag for concurrent CDP crawls When True, forces creation of a new browser context instead of reusing the default context. Essential for concurrent crawls on the same browser to prevent navigation conflicts. * Add context caching to create_isolated_context branch Uses contexts_by_config cache (same as non-CDP mode) to reuse contexts for multiple URLs with same config. Still creates new page per crawl for navigation isolation. Benefits batch/deep crawls. * Add init_scripts support to BrowserConfig for pre-page-load JS injection This adds the ability to inject JavaScript that runs before any page loads, useful for stealth evasions (canvas/audio fingerprinting, userAgentData). - Add init_scripts parameter to BrowserConfig (list of JS strings) - Apply init_scripts in setup_context() via context.add_init_script() - Update from_kwargs() and to_dict() for serialization * Fix CDP connection handling: support WS URLs and proper cleanup Changes to browser_manager.py: 1. _verify_cdp_ready(): Support multiple URL formats - WebSocket URLs (ws://, wss://): Skip HTTP verification, Playwright handles directly - HTTP URLs with query params: Properly parse with urlparse to preserve query string - Fixes issue where naive f"{cdp_url}/json/version" broke WS URLs and query params 2. close(): Proper cleanup when cdp_cleanup_on_close=True - Close all sessions (pages) - Close all contexts - Call browser.close() to disconnect (doesn't terminate browser, just releases connection) - Wait 1 second for CDP connection to fully release - Stop Playwright instance to prevent memory leaks This enables: - Connecting to specific browsers via WS URL - Reusing the same browser with multiple sequential connections - No user wait needed between connections (internal 1s delay handles it) Added tests/browser/test_cdp_cleanup_reuse.py with comprehensive tests. * Update gitignore * Some debugging for caching * Add _generate_screenshot_from_html for raw: and file:// URLs Implements the missing method that was being called but never defined. Now raw: and file:// URLs can generate screenshots by: 1. Loading HTML into a browser page via page.set_content() 2. Taking screenshot using existing take_screenshot() method 3. Cleaning up the page afterward This enables cached HTML to be rendered with screenshots in crawl4ai-cloud. * Add PDF and MHTML support for raw: and file:// URLs - Replace _generate_screenshot_from_html with _generate_media_from_html - New method handles screenshot, PDF, and MHTML in one browser session - Update raw: and file:// URL handlers to use new method - Enables cached HTML to generate all media types * Add crash recovery for deep crawl strategies Add optional resume_state and on_state_change parameters to all deep crawl strategies (BFS, DFS, Best-First) for cloud deployment crash recovery. Features: - resume_state: Pass saved state to resume from checkpoint - on_state_change: Async callback fired after each URL for real-time state persistence to external storage (Redis, DB, etc.) - export_state(): Get last captured state manually - Zero overhead when features are disabled (None defaults) State includes visited URLs, pending queue/stack, depths, and pages_crawled count. All state is JSON-serializable. * Fix: HTTP strategy raw: URL parsing truncates at # character The AsyncHTTPCrawlerStrategy.crawl() method used urlparse() to extract content from raw: URLs. This caused HTML with CSS color codes like #eee to be truncated because # is treated as a URL fragment delimiter. Before: raw:body{background:#eee} -> parsed.path = 'body{background:' After: raw:body{background:#eee} -> raw_content = 'body{background:#eee' Fix: Strip the raw: or raw:// prefix directly instead of using urlparse, matching how the browser strategy handles it. * Add base_url parameter to CrawlerRunConfig for raw HTML processing When processing raw: HTML (e.g., from cache), the URL parameter is meaningless for markdown link resolution. This adds a base_url parameter that can be set explicitly to provide proper URL resolution context. Changes: - Add base_url parameter to CrawlerRunConfig.__init__ - Add base_url to CrawlerRunConfig.from_kwargs - Update aprocess_html to use base_url for markdown generation Usage: config = CrawlerRunConfig(base_url='https://example.com') result = await crawler.arun(url='raw:{html}', config=config) * Add prefetch mode for two-phase deep crawling - Add `prefetch` parameter to CrawlerRunConfig - Add `quick_extract_links()` function for fast link extraction - Add short-circuit in aprocess_html() for prefetch mode - Add 42 tests (unit, integration, regression) 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com> * Updates on proxy rotation and proxy configuration * Add proxy support to HTTP crawler strategy * Add browser pipeline support for raw:/file:// URLs - Add process_in_browser parameter to CrawlerRunConfig - Route raw:/file:// URLs through _crawl_web() when browser operations needed - Use page.set_content() instead of goto() for local content - Fix cookie handling for non-HTTP URLs in browser_manager - Auto-detect browser requirements: js_code, wait_for, screenshot, etc. - Maintain fast path for raw:/file:// without browser params Fixes #310 * Add smart TTL cache for sitemap URL seeder - Add cache_ttl_hours and validate_sitemap_lastmod params to SeedingConfig - New JSON cache format with metadata (version, created_at, lastmod, url_count) - Cache validation by TTL expiry and sitemap lastmod comparison - Auto-migration from old .jsonl to new .json format - Fixes bug where incomplete cache was used indefinitely * Update URL seeder docs with smart TTL cache parameters - Add cache_ttl_hours and validate_sitemap_lastmod to parameter table - Document smart TTL cache validation with examples - Add cache-related troubleshooting entries - Update key features summary * Add MEMORY.md to gitignore * Docs: Add multi-sample schema generation section Add documentation explaining how to pass multiple HTML samples to generate_schema() for stable selectors that work across pages with varying DOM structures. Includes: - Problem explanation (fragile nth-child selectors) - Solution with code example - Key points for multi-sample queries - Comparison table of fragile vs stable selectors * Fix critical RCE and LFI vulnerabilities in Docker API deployment Security fixes for vulnerabilities reported by ProjectDiscovery: 1. Remote Code Execution via Hooks (CVE pending) - Remove __import__ from allowed_builtins in hook_manager.py - Prevents arbitrary module imports (os, subprocess, etc.) - Hooks now disabled by default via CRAWL4AI_HOOKS_ENABLED env var 2. Local File Inclusion via file:// URLs (CVE pending) - Add URL scheme validation to /execute_js, /screenshot, /pdf, /html - Block file://, javascript:, data: and other dangerous schemes - Only allow http://, https://, and raw: (where appropriate) 3. Security hardening - Add CRAWL4AI_HOOKS_ENABLED=false as default (opt-in for hooks) - Add security warning comments in config.yml - Add validate_url_scheme() helper for consistent validation Testing: - Add unit tests (test_security_fixes.py) - 16 tests - Add integration tests (run_security_tests.py) for live server Affected endpoints: - POST /crawl (hooks disabled by default) - POST /crawl/stream (hooks disabled by default) - POST /execute_js (URL validation added) - POST /screenshot (URL validation added) - POST /pdf (URL validation added) - POST /html (URL validation added) Breaking changes: - Hooks require CRAWL4AI_HOOKS_ENABLED=true to function - file:// URLs no longer work on API endpoints (use library directly) * Enhance authentication flow by implementing JWT token retrieval and adding authorization headers to API requests * Add release notes for v0.7.9, detailing breaking changes, security fixes, new features, bug fixes, and documentation updates * Add release notes for v0.8.0, detailing breaking changes, security fixes, new features, bug fixes, and documentation updates Documentation for v0.8.0 release: - SECURITY.md: Security policy and vulnerability reporting guidelines - RELEASE_NOTES_v0.8.0.md: Comprehensive release notes - migration/v0.8.0-upgrade-guide.md: Step-by-step migration guide - security/GHSA-DRAFT-RCE-LFI.md: GitHub security advisory drafts - CHANGELOG.md: Updated with v0.8.0 changes Breaking changes documented: - Docker API hooks disabled by default (CRAWL4AI_HOOKS_ENABLED) - file:// URLs blocked on Docker API endpoints Security fixes credited to Neo by ProjectDiscovery * Add examples for deep crawl crash recovery and prefetch mode in documentation * Release v0.8.0: The v0.8.0 Update - Updated version to 0.8.0 - Added comprehensive demo and release notes - Updated all documentation * Update security researcher acknowledgment with a hyperlink for Neo by ProjectDiscovery * Add async agenerate_schema method for schema generation - Extract prompt building to shared _build_schema_prompt() method - Add agenerate_schema() async version using aperform_completion_with_backoff - Refactor generate_schema() to use shared prompt builder - Fixes Gemini/Vertex AI compatibility in async contexts (FastAPI) * Fix: Enable litellm.drop_params for O-series/GPT-5 model compatibility O-series (o1, o3) and GPT-5 models only support temperature=1. Setting litellm.drop_params=True auto-drops unsupported parameters instead of throwing UnsupportedParamsError. Fixes temperature=0.01 error for these models in LLM extraction. --------- Co-authored-by: rbushria <rbushri@gmail.com> Co-authored-by: AHMET YILMAZ <tawfik@kidocode.com> Co-authored-by: Soham Kukreti <kukretisoham@gmail.com> Co-authored-by: Chris Murphy <chris.murphy@klaviyo.com> Co-authored-by: unclecode <unclecode@kidocode.com> Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
679 lines
26 KiB
Python
679 lines
26 KiB
Python
import os
|
|
import time
|
|
from pathlib import Path
|
|
import aiosqlite
|
|
import asyncio
|
|
from typing import Optional, Dict
|
|
from contextlib import asynccontextmanager
|
|
import json
|
|
from .models import CrawlResult, MarkdownGenerationResult, StringCompatibleMarkdown
|
|
import aiofiles
|
|
from .async_logger import AsyncLogger
|
|
|
|
from .utils import ensure_content_dirs, generate_content_hash
|
|
from .utils import VersionManager
|
|
from .utils import get_error_context, create_box_message
|
|
|
|
base_directory = DB_PATH = os.path.join(
|
|
os.getenv("CRAWL4_AI_BASE_DIRECTORY", Path.home()), ".crawl4ai"
|
|
)
|
|
os.makedirs(DB_PATH, exist_ok=True)
|
|
DB_PATH = os.path.join(base_directory, "crawl4ai.db")
|
|
|
|
|
|
class AsyncDatabaseManager:
|
|
def __init__(self, pool_size: int = 10, max_retries: int = 3):
|
|
self.db_path = DB_PATH
|
|
self.content_paths = ensure_content_dirs(os.path.dirname(DB_PATH))
|
|
self.pool_size = pool_size
|
|
self.max_retries = max_retries
|
|
self.connection_pool: Dict[int, aiosqlite.Connection] = {}
|
|
self.pool_lock = asyncio.Lock()
|
|
self.init_lock = asyncio.Lock()
|
|
self.connection_semaphore = asyncio.Semaphore(pool_size)
|
|
self._initialized = False
|
|
self.version_manager = VersionManager()
|
|
self.logger = AsyncLogger(
|
|
log_file=os.path.join(base_directory, ".crawl4ai", "crawler_db.log"),
|
|
verbose=False,
|
|
tag_width=10,
|
|
)
|
|
|
|
async def initialize(self):
|
|
"""Initialize the database and connection pool"""
|
|
try:
|
|
self.logger.info("Initializing database", tag="INIT")
|
|
# Ensure the database file exists
|
|
os.makedirs(os.path.dirname(self.db_path), exist_ok=True)
|
|
|
|
# Check if version update is needed
|
|
needs_update = self.version_manager.needs_update()
|
|
|
|
# Always ensure base table exists
|
|
await self.ainit_db()
|
|
|
|
# Verify the table exists
|
|
async with aiosqlite.connect(self.db_path, timeout=30.0) as db:
|
|
async with db.execute(
|
|
"SELECT name FROM sqlite_master WHERE type='table' AND name='crawled_data'"
|
|
) as cursor:
|
|
result = await cursor.fetchone()
|
|
if not result:
|
|
raise Exception("crawled_data table was not created")
|
|
|
|
# If version changed or fresh install, run updates
|
|
if needs_update:
|
|
self.logger.info("New version detected, running updates", tag="INIT")
|
|
await self.update_db_schema()
|
|
from .migrations import (
|
|
run_migration,
|
|
) # Import here to avoid circular imports
|
|
|
|
await run_migration()
|
|
self.version_manager.update_version() # Update stored version after successful migration
|
|
self.logger.success(
|
|
"Version update completed successfully", tag="COMPLETE"
|
|
)
|
|
else:
|
|
self.logger.success(
|
|
"Database initialization completed successfully", tag="COMPLETE"
|
|
)
|
|
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Database initialization error: {error}",
|
|
tag="ERROR",
|
|
params={"error": str(e)},
|
|
)
|
|
self.logger.info(
|
|
message="Database will be initialized on first use", tag="INIT"
|
|
)
|
|
|
|
raise
|
|
|
|
async def cleanup(self):
|
|
"""Cleanup connections when shutting down"""
|
|
async with self.pool_lock:
|
|
for conn in self.connection_pool.values():
|
|
await conn.close()
|
|
self.connection_pool.clear()
|
|
|
|
@asynccontextmanager
|
|
async def get_connection(self):
|
|
"""Connection pool manager with enhanced error handling"""
|
|
if not self._initialized:
|
|
async with self.init_lock:
|
|
if not self._initialized:
|
|
try:
|
|
await self.initialize()
|
|
self._initialized = True
|
|
except Exception as e:
|
|
import sys
|
|
|
|
error_context = get_error_context(sys.exc_info())
|
|
self.logger.error(
|
|
message="Database initialization failed:\n{error}\n\nContext:\n{context}\n\nTraceback:\n{traceback}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={
|
|
"error": str(e),
|
|
"context": error_context["code_context"],
|
|
"traceback": error_context["full_traceback"],
|
|
},
|
|
)
|
|
raise
|
|
|
|
await self.connection_semaphore.acquire()
|
|
task_id = id(asyncio.current_task())
|
|
|
|
try:
|
|
async with self.pool_lock:
|
|
if task_id not in self.connection_pool:
|
|
try:
|
|
conn = await aiosqlite.connect(self.db_path, timeout=30.0)
|
|
await conn.execute("PRAGMA journal_mode = WAL")
|
|
await conn.execute("PRAGMA busy_timeout = 5000")
|
|
|
|
# Verify database structure
|
|
async with conn.execute(
|
|
"PRAGMA table_info(crawled_data)"
|
|
) as cursor:
|
|
columns = await cursor.fetchall()
|
|
column_names = [col[1] for col in columns]
|
|
expected_columns = {
|
|
"url",
|
|
"html",
|
|
"cleaned_html",
|
|
"markdown",
|
|
"extracted_content",
|
|
"success",
|
|
"media",
|
|
"links",
|
|
"metadata",
|
|
"screenshot",
|
|
"response_headers",
|
|
"downloaded_files",
|
|
}
|
|
missing_columns = expected_columns - set(column_names)
|
|
if missing_columns:
|
|
raise ValueError(
|
|
f"Database missing columns: {missing_columns}"
|
|
)
|
|
|
|
self.connection_pool[task_id] = conn
|
|
except Exception as e:
|
|
import sys
|
|
|
|
error_context = get_error_context(sys.exc_info())
|
|
error_message = (
|
|
f"Unexpected error in db get_connection at line {error_context['line_no']} "
|
|
f"in {error_context['function']} ({error_context['filename']}):\n"
|
|
f"Error: {str(e)}\n\n"
|
|
f"Code context:\n{error_context['code_context']}"
|
|
)
|
|
self.logger.error(
|
|
message="{error}",
|
|
tag="ERROR",
|
|
params={"error": str(error_message)},
|
|
boxes=["error"],
|
|
)
|
|
|
|
raise
|
|
|
|
yield self.connection_pool[task_id]
|
|
|
|
except Exception as e:
|
|
import sys
|
|
|
|
error_context = get_error_context(sys.exc_info())
|
|
error_message = (
|
|
f"Unexpected error in db get_connection at line {error_context['line_no']} "
|
|
f"in {error_context['function']} ({error_context['filename']}):\n"
|
|
f"Error: {str(e)}\n\n"
|
|
f"Code context:\n{error_context['code_context']}"
|
|
)
|
|
self.logger.error(
|
|
message="{error}",
|
|
tag="ERROR",
|
|
params={"error": str(error_message)},
|
|
boxes=["error"],
|
|
)
|
|
raise
|
|
finally:
|
|
async with self.pool_lock:
|
|
if task_id in self.connection_pool:
|
|
await self.connection_pool[task_id].close()
|
|
del self.connection_pool[task_id]
|
|
self.connection_semaphore.release()
|
|
|
|
async def execute_with_retry(self, operation, *args):
|
|
"""Execute database operations with retry logic"""
|
|
for attempt in range(self.max_retries):
|
|
try:
|
|
async with self.get_connection() as db:
|
|
result = await operation(db, *args)
|
|
await db.commit()
|
|
return result
|
|
except Exception as e:
|
|
if attempt == self.max_retries - 1:
|
|
self.logger.error(
|
|
message="Operation failed after {retries} attempts: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"retries": self.max_retries, "error": str(e)},
|
|
)
|
|
raise
|
|
await asyncio.sleep(1 * (attempt + 1)) # Exponential backoff
|
|
|
|
async def ainit_db(self):
|
|
"""Initialize database schema"""
|
|
async with aiosqlite.connect(self.db_path, timeout=30.0) as db:
|
|
await db.execute(
|
|
"""
|
|
CREATE TABLE IF NOT EXISTS crawled_data (
|
|
url TEXT PRIMARY KEY,
|
|
html TEXT,
|
|
cleaned_html TEXT,
|
|
markdown TEXT,
|
|
extracted_content TEXT,
|
|
success BOOLEAN,
|
|
media TEXT DEFAULT "{}",
|
|
links TEXT DEFAULT "{}",
|
|
metadata TEXT DEFAULT "{}",
|
|
screenshot TEXT DEFAULT "",
|
|
response_headers TEXT DEFAULT "{}",
|
|
downloaded_files TEXT DEFAULT "{}" -- New column added
|
|
)
|
|
"""
|
|
)
|
|
await db.commit()
|
|
|
|
async def update_db_schema(self):
|
|
"""Update database schema if needed"""
|
|
async with aiosqlite.connect(self.db_path, timeout=30.0) as db:
|
|
cursor = await db.execute("PRAGMA table_info(crawled_data)")
|
|
columns = await cursor.fetchall()
|
|
column_names = [column[1] for column in columns]
|
|
|
|
# List of new columns to add
|
|
new_columns = [
|
|
"media",
|
|
"links",
|
|
"metadata",
|
|
"screenshot",
|
|
"response_headers",
|
|
"downloaded_files",
|
|
# Smart cache validation columns (added in 0.8.x)
|
|
"etag",
|
|
"last_modified",
|
|
"head_fingerprint",
|
|
"cached_at",
|
|
]
|
|
|
|
for column in new_columns:
|
|
if column not in column_names:
|
|
await self.aalter_db_add_column(column, db)
|
|
await db.commit()
|
|
|
|
async def aalter_db_add_column(self, new_column: str, db):
|
|
"""Add new column to the database"""
|
|
if new_column == "response_headers":
|
|
await db.execute(
|
|
f'ALTER TABLE crawled_data ADD COLUMN {new_column} TEXT DEFAULT "{{}}"'
|
|
)
|
|
elif new_column == "cached_at":
|
|
# Timestamp column for cache validation
|
|
await db.execute(
|
|
f"ALTER TABLE crawled_data ADD COLUMN {new_column} REAL DEFAULT 0"
|
|
)
|
|
else:
|
|
await db.execute(
|
|
f'ALTER TABLE crawled_data ADD COLUMN {new_column} TEXT DEFAULT ""'
|
|
)
|
|
self.logger.info(
|
|
message="Added column '{column}' to the database",
|
|
tag="INIT",
|
|
params={"column": new_column},
|
|
)
|
|
|
|
async def aget_cached_url(self, url: str) -> Optional[CrawlResult]:
|
|
"""Retrieve cached URL data as CrawlResult"""
|
|
|
|
async def _get(db):
|
|
async with db.execute(
|
|
"SELECT * FROM crawled_data WHERE url = ?", (url,)
|
|
) as cursor:
|
|
row = await cursor.fetchone()
|
|
if not row:
|
|
return None
|
|
|
|
# Get column names
|
|
columns = [description[0] for description in cursor.description]
|
|
# Create dict from row data
|
|
row_dict = dict(zip(columns, row))
|
|
|
|
# Load content from files using stored hashes
|
|
content_fields = {
|
|
"html": row_dict["html"],
|
|
"cleaned_html": row_dict["cleaned_html"],
|
|
"markdown": row_dict["markdown"],
|
|
"extracted_content": row_dict["extracted_content"],
|
|
"screenshot": row_dict["screenshot"],
|
|
"screenshots": row_dict["screenshot"],
|
|
}
|
|
|
|
for field, hash_value in content_fields.items():
|
|
if hash_value:
|
|
content = await self._load_content(
|
|
hash_value,
|
|
field.split("_")[0], # Get content type from field name
|
|
)
|
|
row_dict[field] = content or ""
|
|
else:
|
|
row_dict[field] = ""
|
|
|
|
# Parse JSON fields
|
|
json_fields = [
|
|
"media",
|
|
"links",
|
|
"metadata",
|
|
"response_headers",
|
|
"markdown",
|
|
]
|
|
for field in json_fields:
|
|
try:
|
|
row_dict[field] = (
|
|
json.loads(row_dict[field]) if row_dict[field] else {}
|
|
)
|
|
except json.JSONDecodeError:
|
|
# Very UGLY, never mention it to me please
|
|
if field == "markdown" and isinstance(row_dict[field], str):
|
|
row_dict[field] = MarkdownGenerationResult(
|
|
raw_markdown=row_dict[field] or "",
|
|
markdown_with_citations="",
|
|
references_markdown="",
|
|
fit_markdown="",
|
|
fit_html="",
|
|
)
|
|
else:
|
|
row_dict[field] = {}
|
|
|
|
if isinstance(row_dict["markdown"], Dict):
|
|
if row_dict["markdown"].get("raw_markdown"):
|
|
row_dict["markdown"] = row_dict["markdown"]["raw_markdown"]
|
|
|
|
# Parse downloaded_files
|
|
try:
|
|
row_dict["downloaded_files"] = (
|
|
json.loads(row_dict["downloaded_files"])
|
|
if row_dict["downloaded_files"]
|
|
else []
|
|
)
|
|
except json.JSONDecodeError:
|
|
row_dict["downloaded_files"] = []
|
|
|
|
# Remove any fields not in CrawlResult model
|
|
valid_fields = CrawlResult.__annotations__.keys()
|
|
filtered_dict = {k: v for k, v in row_dict.items() if k in valid_fields}
|
|
filtered_dict["markdown"] = row_dict["markdown"]
|
|
return CrawlResult(**filtered_dict)
|
|
|
|
try:
|
|
return await self.execute_with_retry(_get)
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Error retrieving cached URL: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"error": str(e)},
|
|
)
|
|
return None
|
|
|
|
async def aget_cache_metadata(self, url: str) -> Optional[Dict]:
|
|
"""
|
|
Retrieve only cache validation metadata for a URL (lightweight query).
|
|
|
|
Returns dict with: url, etag, last_modified, head_fingerprint, cached_at, response_headers
|
|
This is used for cache validation without loading full content.
|
|
"""
|
|
async def _get_metadata(db):
|
|
async with db.execute(
|
|
"""SELECT url, etag, last_modified, head_fingerprint, cached_at, response_headers
|
|
FROM crawled_data WHERE url = ?""",
|
|
(url,)
|
|
) as cursor:
|
|
row = await cursor.fetchone()
|
|
if not row:
|
|
return None
|
|
|
|
columns = [description[0] for description in cursor.description]
|
|
row_dict = dict(zip(columns, row))
|
|
|
|
# Parse response_headers JSON
|
|
try:
|
|
row_dict["response_headers"] = (
|
|
json.loads(row_dict["response_headers"])
|
|
if row_dict["response_headers"] else {}
|
|
)
|
|
except json.JSONDecodeError:
|
|
row_dict["response_headers"] = {}
|
|
|
|
return row_dict
|
|
|
|
try:
|
|
return await self.execute_with_retry(_get_metadata)
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Error retrieving cache metadata: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"error": str(e)},
|
|
)
|
|
return None
|
|
|
|
async def aupdate_cache_metadata(
|
|
self,
|
|
url: str,
|
|
etag: Optional[str] = None,
|
|
last_modified: Optional[str] = None,
|
|
head_fingerprint: Optional[str] = None,
|
|
):
|
|
"""
|
|
Update only the cache validation metadata for a URL.
|
|
Used to update etag/last_modified after a successful validation.
|
|
"""
|
|
async def _update(db):
|
|
updates = []
|
|
values = []
|
|
|
|
if etag is not None:
|
|
updates.append("etag = ?")
|
|
values.append(etag)
|
|
if last_modified is not None:
|
|
updates.append("last_modified = ?")
|
|
values.append(last_modified)
|
|
if head_fingerprint is not None:
|
|
updates.append("head_fingerprint = ?")
|
|
values.append(head_fingerprint)
|
|
|
|
if not updates:
|
|
return
|
|
|
|
values.append(url)
|
|
await db.execute(
|
|
f"UPDATE crawled_data SET {', '.join(updates)} WHERE url = ?",
|
|
tuple(values)
|
|
)
|
|
|
|
try:
|
|
await self.execute_with_retry(_update)
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Error updating cache metadata: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"error": str(e)},
|
|
)
|
|
|
|
async def acache_url(self, result: CrawlResult):
|
|
"""Cache CrawlResult data"""
|
|
# Store content files and get hashes
|
|
content_map = {
|
|
"html": (result.html, "html"),
|
|
"cleaned_html": (result.cleaned_html or "", "cleaned"),
|
|
"markdown": None,
|
|
"extracted_content": (result.extracted_content or "", "extracted"),
|
|
"screenshot": (result.screenshot or "", "screenshots"),
|
|
}
|
|
|
|
try:
|
|
if isinstance(result.markdown, StringCompatibleMarkdown):
|
|
content_map["markdown"] = (
|
|
result.markdown,
|
|
"markdown",
|
|
)
|
|
elif isinstance(result.markdown, MarkdownGenerationResult):
|
|
content_map["markdown"] = (
|
|
result.markdown.model_dump_json(),
|
|
"markdown",
|
|
)
|
|
elif isinstance(result.markdown, str):
|
|
markdown_result = MarkdownGenerationResult(raw_markdown=result.markdown)
|
|
content_map["markdown"] = (
|
|
markdown_result.model_dump_json(),
|
|
"markdown",
|
|
)
|
|
else:
|
|
content_map["markdown"] = (
|
|
MarkdownGenerationResult().model_dump_json(),
|
|
"markdown",
|
|
)
|
|
except Exception as e:
|
|
self.logger.warning(
|
|
message=f"Error processing markdown content: {str(e)}", tag="WARNING"
|
|
)
|
|
# Fallback to empty markdown result
|
|
content_map["markdown"] = (
|
|
MarkdownGenerationResult().model_dump_json(),
|
|
"markdown",
|
|
)
|
|
|
|
content_hashes = {}
|
|
for field, (content, content_type) in content_map.items():
|
|
content_hashes[field] = await self._store_content(content, content_type)
|
|
|
|
# Extract cache validation headers from response
|
|
response_headers = result.response_headers or {}
|
|
etag = response_headers.get("etag") or response_headers.get("ETag") or ""
|
|
last_modified = response_headers.get("last-modified") or response_headers.get("Last-Modified") or ""
|
|
# head_fingerprint is set by caller via result attribute (if available)
|
|
head_fingerprint = getattr(result, "head_fingerprint", None) or ""
|
|
cached_at = time.time()
|
|
|
|
async def _cache(db):
|
|
await db.execute(
|
|
"""
|
|
INSERT INTO crawled_data (
|
|
url, html, cleaned_html, markdown,
|
|
extracted_content, success, media, links, metadata,
|
|
screenshot, response_headers, downloaded_files,
|
|
etag, last_modified, head_fingerprint, cached_at
|
|
)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT(url) DO UPDATE SET
|
|
html = excluded.html,
|
|
cleaned_html = excluded.cleaned_html,
|
|
markdown = excluded.markdown,
|
|
extracted_content = excluded.extracted_content,
|
|
success = excluded.success,
|
|
media = excluded.media,
|
|
links = excluded.links,
|
|
metadata = excluded.metadata,
|
|
screenshot = excluded.screenshot,
|
|
response_headers = excluded.response_headers,
|
|
downloaded_files = excluded.downloaded_files,
|
|
etag = excluded.etag,
|
|
last_modified = excluded.last_modified,
|
|
head_fingerprint = excluded.head_fingerprint,
|
|
cached_at = excluded.cached_at
|
|
""",
|
|
(
|
|
result.url,
|
|
content_hashes["html"],
|
|
content_hashes["cleaned_html"],
|
|
content_hashes["markdown"],
|
|
content_hashes["extracted_content"],
|
|
result.success,
|
|
json.dumps(result.media),
|
|
json.dumps(result.links),
|
|
json.dumps(result.metadata or {}),
|
|
content_hashes["screenshot"],
|
|
json.dumps(result.response_headers or {}),
|
|
json.dumps(result.downloaded_files or []),
|
|
etag,
|
|
last_modified,
|
|
head_fingerprint,
|
|
cached_at,
|
|
),
|
|
)
|
|
|
|
try:
|
|
await self.execute_with_retry(_cache)
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Error caching URL: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"error": str(e)},
|
|
)
|
|
|
|
async def aget_total_count(self) -> int:
|
|
"""Get total number of cached URLs"""
|
|
|
|
async def _count(db):
|
|
async with db.execute("SELECT COUNT(*) FROM crawled_data") as cursor:
|
|
result = await cursor.fetchone()
|
|
return result[0] if result else 0
|
|
|
|
try:
|
|
return await self.execute_with_retry(_count)
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Error getting total count: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"error": str(e)},
|
|
)
|
|
return 0
|
|
|
|
async def aclear_db(self):
|
|
"""Clear all data from the database"""
|
|
|
|
async def _clear(db):
|
|
await db.execute("DELETE FROM crawled_data")
|
|
|
|
try:
|
|
await self.execute_with_retry(_clear)
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Error clearing database: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"error": str(e)},
|
|
)
|
|
|
|
async def aflush_db(self):
|
|
"""Drop the entire table"""
|
|
|
|
async def _flush(db):
|
|
await db.execute("DROP TABLE IF EXISTS crawled_data")
|
|
|
|
try:
|
|
await self.execute_with_retry(_flush)
|
|
except Exception as e:
|
|
self.logger.error(
|
|
message="Error flushing database: {error}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"error": str(e)},
|
|
)
|
|
|
|
async def _store_content(self, content: str, content_type: str) -> str:
|
|
"""Store content in filesystem and return hash"""
|
|
if not content:
|
|
return ""
|
|
|
|
content_hash = generate_content_hash(content)
|
|
file_path = os.path.join(self.content_paths[content_type], content_hash)
|
|
|
|
# Only write if file doesn't exist
|
|
if not os.path.exists(file_path):
|
|
async with aiofiles.open(file_path, "w", encoding="utf-8") as f:
|
|
await f.write(content)
|
|
|
|
return content_hash
|
|
|
|
async def _load_content(
|
|
self, content_hash: str, content_type: str
|
|
) -> Optional[str]:
|
|
"""Load content from filesystem by hash"""
|
|
if not content_hash:
|
|
return None
|
|
|
|
file_path = os.path.join(self.content_paths[content_type], content_hash)
|
|
try:
|
|
async with aiofiles.open(file_path, "r", encoding="utf-8") as f:
|
|
return await f.read()
|
|
except:
|
|
self.logger.error(
|
|
message="Failed to load content: {file_path}",
|
|
tag="ERROR",
|
|
force_verbose=True,
|
|
params={"file_path": file_path},
|
|
)
|
|
return None
|
|
|
|
|
|
# Create a singleton instance
|
|
async_db_manager = AsyncDatabaseManager()
|