pig.py/neopig.py
Russell Ballestrini bc11c516b6 Add filevault system with async wrapper and convert sync ops to async
New files:
- filevault.py: Hash-based file storage with use_pairs option (v1.1.0)
- async_filevault.py: Async wrapper using asyncio.to_thread()
- domain_vault.py: Triple vault system for web archival (HTML, Media, Linkpeek)
- screenshot.py: Async screenshot capture using uri2png
- tests/unit/: Comprehensive test suite (85 tests)

Sync-to-async conversions:
- storage.py: Wrap Path operations in asyncio.to_thread()
- domain_vault.py: Wrap exists(), mkdir(), rglob(), os.walk() in asyncio.to_thread()
- screenshot.py: Wrap read_bytes(), write_bytes(), unlink() in asyncio.to_thread()

All sync filesystem operations now run in thread pool to avoid blocking async loop.
2025-12-22 20:18:20 -05:00

631 lines
22 KiB
Python

#!/usr/bin/env python3
"""
neopig - Neo Python Image Grabber
Full domain crawler that:
1. Crawls entire domain (no depth limit)
2. Extracts all images/videos via multiple methods
3. Uses HEAD checks for extensionless URLs
4. Stores in vault (MD5 dedupe)
5. Queues for Qwen 3 VL analysis
6. Indexes metadata in SQLite
Based on pig.py by Russell Ballestrini
https://russell.ballestrini.net/python-image-grabber-pig-py/
Usage:
python neopig.py https://example.com "rick and morty" --mode images
python neopig.py https://example.com --mode all --depth -1 # Full domain slurp
"""
import argparse
import asyncio
import hashlib
import logging
import sys
from datetime import datetime, timezone
from pathlib import Path
from typing import List, Dict, Any, Optional, Set
from urllib.parse import urlparse
from async_web_fetcher import (
AsyncWebFetcher,
CrawlMode,
MediaItem,
extract_media_from_html,
get_media_type_from_extension,
get_media_type_from_mime,
)
from storage import ImageVault
from database import Database
from screenshot import ScreenshotCapture, ScreenshotConfig
from domain_vault import VaultManager, DomainHtmlVault, DomainMediaVault, DomainLinkpeekVault, extract_media_urls
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)
class NeoPig:
"""
Neo Python Image Grabber - async media crawler with deduplication.
"""
def __init__(
self,
db_path: str = "neopig.db",
vault_path: str = "vault",
user_agent: str = "neopig/1.0 (ethical image crawler)",
screenshot_config: ScreenshotConfig = None,
):
self.db = Database(db_path)
self.vault = ImageVault(vault_path)
# Triple filevault system: html_vault/, media_vault/, and linkpeek_vault/
self.domain_vaults = VaultManager(
html_vault_base=f"{vault_path}/html_vault",
media_vault_base=f"{vault_path}/media_vault",
linkpeek_vault_base=f"{vault_path}/linkpeek_vault",
media_base_url='/media',
linkpeek_base_url='/linkpeek',
)
self.fetcher = AsyncWebFetcher(user_agent=user_agent)
self.screenshot = ScreenshotCapture(screenshot_config or ScreenshotConfig())
self.screenshot_config = screenshot_config or ScreenshotConfig()
self.vault_path = vault_path
# Track stats
self.stats = {
'pages_crawled': 0,
'pages_changed': 0, # Pages with content changes (for git commit)
'media_found': 0,
'media_downloaded': 0,
'media_new': 0, # New media (for git commit)
'duplicates_skipped': 0,
'screenshots_taken': 0,
'errors': 0,
}
# Track seen media URLs to avoid re-processing
self.seen_media: Set[str] = set()
# Track screenshotted pages to avoid duplicates
self.seen_screenshots: Set[str] = set()
# Track per-domain stats for vault commits
self._domain_stats: Dict[str, Dict[str, int]] = {} # domain -> {pages_changed, media_new, screenshots_new}
def _get_domain(self, url: str) -> str:
"""Extract domain from URL."""
parsed = urlparse(url)
return parsed.netloc.lower()
def _track_domain_stat(self, domain: str, stat: str, increment: int = 1):
"""Track per-domain stats for vault commits."""
if domain not in self._domain_stats:
self._domain_stats[domain] = {'pages_changed': 0, 'media_new': 0, 'screenshots_new': 0}
self._domain_stats[domain][stat] += increment
async def _archive_page_to_vault(
self,
url: str,
html: str,
media_mappings: Dict[str, str] = None,
):
"""Archive a page to the HTML vault."""
domain = self._get_domain(url)
html_vault = self.domain_vaults.get_html_vault(domain)
is_changed, _ = await html_vault.archive_page(url, html, media_mappings)
if is_changed:
self._track_domain_stat(domain, 'pages_changed')
self.stats['pages_changed'] += 1
async def _archive_media_to_vault(
self,
url: str,
content: bytes,
page_url: str = '',
):
"""Archive media to the media vault."""
domain = self._get_domain(url)
media_vault = self.domain_vaults.get_media_vault(domain)
is_new, _, _ = await media_vault.archive_media(url, content, page_url)
if is_new:
self._track_domain_stat(domain, 'media_new')
self.stats['media_new'] += 1
async def _archive_screenshot_to_vault(
self,
url: str,
screenshot_data: bytes,
):
"""Archive screenshot to the linkpeek vault."""
domain = self._get_domain(url)
linkpeek_vault = self.domain_vaults.get_linkpeek_vault(domain)
is_new, _, _ = await linkpeek_vault.archive_screenshot(url, screenshot_data)
if is_new:
self._track_domain_stat(domain, 'screenshots_new')
async def _finish_domain_vaults(self, keywords: List[str] = None):
"""Commit changes to all domain vaults that have diffs."""
for domain, stats in self._domain_stats.items():
if stats['pages_changed'] > 0 or stats['media_new'] > 0 or stats['screenshots_new'] > 0:
html_vault = self.domain_vaults.get_html_vault(domain)
media_vault = self.domain_vaults.get_media_vault(domain)
linkpeek_vault = self.domain_vaults.get_linkpeek_vault(domain)
html_commit = await html_vault.finish_crawl(stats)
media_commit = await media_vault.finish_crawl(stats)
linkpeek_commit = await linkpeek_vault.finish_crawl(stats)
if html_commit:
logger.info(f"Committed HTML vault for {domain}: {html_commit[:8]}")
if media_commit:
logger.info(f"Committed media vault for {domain}: {media_commit[:8]}")
if linkpeek_commit:
logger.info(f"Committed linkpeek vault for {domain}: {linkpeek_commit[:8]}")
async def init(self):
"""Initialize database and vault."""
await self.db.init()
await self.vault.init()
async def crawl(
self,
target_uri: str,
keywords: List[str] = None,
mode: CrawlMode = CrawlMode.IMAGES,
depth: int = -1, # -1 = unlimited
max_pages: int = -1, # -1 = unlimited
download_media: bool = True,
) -> Dict[str, Any]:
"""
Crawl a domain for media.
Args:
target_uri: Starting URL
keywords: Keywords to tag media with (e.g., ["rick and morty"])
mode: CrawlMode - what to collect (IMAGES, VIDEOS, MEDIA, ALL)
depth: Crawl depth (-1 = unlimited)
max_pages: Max pages to crawl (-1 = unlimited)
download_media: Whether to download media or just index URLs
Returns:
Crawl statistics
"""
keywords = keywords or []
# Create crawl job
job_id = await self.db.create_crawl_job(
target_uri=target_uri,
keywords=keywords,
mode=mode.value,
)
logger.info(f"Starting neopig crawl job {job_id}")
logger.info(f"Target: {target_uri}")
logger.info(f"Mode: {mode.value}")
logger.info(f"Keywords: {keywords}")
logger.info(f"Depth: {'unlimited' if depth == -1 else depth}")
# Track timing for stats
start_time = datetime.now(timezone.utc)
last_stats_time = start_time
stats_running = True
# Background task to emit stats every 15 seconds
async def stats_reporter():
nonlocal last_stats_time
while stats_running:
await asyncio.sleep(15)
if not stats_running:
break
elapsed = (datetime.now(timezone.utc) - start_time).total_seconds()
rate = self.stats['media_downloaded'] / elapsed * 60 if elapsed > 0 else 0
logger.info(f"=== CRAWL STATS ({elapsed:.0f}s) ===")
logger.info(f" Pages: {self.stats['pages_crawled']} | "
f"Found: {self.stats['media_found']} | "
f"Downloaded: {self.stats['media_downloaded']} | "
f"Dupes: {self.stats['duplicates_skipped']} | "
f"Screenshots: {self.stats['screenshots_taken']} | "
f"Errors: {self.stats['errors']}")
logger.info(f" Rate: {rate:.1f}/min | "
f"Vault size: {len(self.seen_media)}")
stats_task = asyncio.create_task(stats_reporter())
# Media callback - called for each discovered media item
async def on_media_discovered(item: Dict[str, Any]):
url = item.get('url')
if not url or url in self.seen_media:
return
self.seen_media.add(url)
self.stats['media_found'] += 1
if download_media:
await self._process_media_item(item, job_id, keywords)
# Note: Screenshots are now captured per-page in on_page_fetched,
# not per-media-item, to honor crawl delay as a unit
# Progress callback
async def on_progress(msg: str):
self.stats['pages_crawled'] += 1
if self.stats['pages_crawled'] % 10 == 0:
logger.info(f"Progress: {self.stats['pages_crawled']} pages, "
f"{self.stats['media_found']} media found, "
f"{self.stats['media_downloaded']} downloaded")
# Page callback - archive raw HTML to vault and capture screenshot
# Screenshot happens here (same crawl delay window as page fetch)
async def on_page_fetched(url: str, html: str):
# For now, archive without media URL rewriting (we'd need to download media first)
# TODO: Build media_mappings after media is downloaded
await self._archive_page_to_vault(url, html, media_mappings=None)
# Capture screenshot for every page (honors crawl delay as a UNIT with page fetch)
if self.screenshot_config.enabled:
await self._capture_page_screenshot(url, job_id, page_title='')
# Run the crawl
pages = await self.fetcher.fetch_with_depth(
start_url=target_uri,
depth=depth,
max_pages=max_pages,
query_keywords=keywords,
mode=mode,
media_callback=on_media_discovered,
progress_callback=on_progress,
page_callback=on_page_fetched,
)
self.stats['pages_crawled'] = len(pages)
# Stop the stats reporter
stats_running = False
stats_task.cancel()
try:
await stats_task
except asyncio.CancelledError:
pass
# Calculate final stats
elapsed = (datetime.now(timezone.utc) - start_time).total_seconds()
rate = self.stats['media_downloaded'] / elapsed * 60 if elapsed > 0 else 0
# Update job status
await self.db.complete_crawl_job(job_id, self.stats)
# Commit domain vaults if there were changes
await self._finish_domain_vaults(keywords)
logger.info(f"=== CRAWL COMPLETE ({elapsed:.0f}s) ===")
logger.info(f" Pages: {self.stats['pages_crawled']} | "
f"Found: {self.stats['media_found']} | "
f"Downloaded: {self.stats['media_downloaded']} | "
f"Dupes: {self.stats['duplicates_skipped']} | "
f"Screenshots: {self.stats['screenshots_taken']} | "
f"Errors: {self.stats['errors']}")
logger.info(f" Rate: {rate:.1f}/min | Total time: {elapsed:.1f}s")
return self.stats
async def _process_media_item(
self,
item: Dict[str, Any],
job_id: int,
keywords: List[str]
):
"""Download and store a media item with page context."""
media_uri = item['url']
page_uri = item.get('source_page', '')
media_type = item.get('media_type', 'unknown')
# Extract page context for searchability
page_title = item.get('page_title', '')
page_description = item.get('page_description', '')
page_keywords = item.get('page_keywords', '')
alt_text = item.get('alt_text', '')
link_text = item.get('link_text', '')
try:
# Check if this exact media+page combo was already crawled
existing_hash = await self.db.check_media_uri_exists(media_uri, page_uri)
if existing_hash:
self.stats['duplicates_skipped'] += 1
logger.debug(f"Already crawled: {media_uri} from {page_uri}")
return
# Check if content already in vault (same MD5 = same content)
# We still need to add the new page context even if content exists
if media_uri in self.seen_media:
# Already processed this media_uri, just add context
# We need to fetch to get MD5, but we can skip if we track it
pass
# Fetch the media
result = await self.fetcher.fetch_media(media_uri)
if not result:
self.stats['errors'] += 1
return
md5_hash = result['md5_hash']
# Check if content already in vault
if await self.vault.exists(md5_hash):
# Content exists, but add this new page context
await self.db.add_media_source(
md5_hash=md5_hash,
media_uri=media_uri,
page_uri=page_uri,
page_title=page_title,
page_description=page_description,
page_keywords=page_keywords,
alt_text=alt_text,
link_text=link_text,
crawl_job_id=job_id,
)
self.stats['duplicates_skipped'] += 1
logger.debug(f"Content exists, added context: {md5_hash} from {page_uri}")
return
# Store in vault (new content)
ext = self._get_extension(media_uri, result.get('mime_type', ''))
await self.vault.store(md5_hash, result['data'], ext)
# Archive to domain media vault (git-tracked)
await self._archive_media_to_vault(media_uri, result['data'], page_uri)
# Record in database with full context
await self.db.create_media_record(
md5_hash=md5_hash,
media_uri=media_uri,
page_uri=page_uri,
crawl_job_id=job_id,
media_type=media_type,
mime_type=result.get('mime_type', ''),
file_size=result.get('size', 0),
page_title=page_title,
page_description=page_description,
page_keywords=page_keywords,
alt_text=alt_text,
link_text=link_text,
)
self.stats['media_downloaded'] += 1
logger.debug(f"Stored: {md5_hash} ({media_uri})")
except Exception as e:
logger.warning(f"Failed to process {media_uri}: {e}")
self.stats['errors'] += 1
async def _capture_page_screenshot(
self,
page_uri: str,
job_id: int,
page_title: str = '',
):
"""Capture and store a screenshot of a page."""
if not self.screenshot_config.enabled:
return
if page_uri in self.seen_screenshots:
return
self.seen_screenshots.add(page_uri)
try:
result = await self.screenshot.capture(page_uri)
if not result:
return
md5_hash = result['md5_hash']
screenshot_data = result['data']
# Store in MD5 vault (for deduplication)
if not await self.vault.exists(md5_hash):
await self.vault.store(md5_hash, screenshot_data, 'png')
# Archive to linkpeek vault (git-tracked by URL path)
await self._archive_screenshot_to_vault(page_uri, screenshot_data)
# Record in database as screenshot type
await self.db.create_media_record(
md5_hash=md5_hash,
media_uri=f"screenshot:{page_uri}",
page_uri=page_uri,
crawl_job_id=job_id,
media_type='screenshot',
mime_type='image/png',
file_size=result.get('size', 0),
page_title=page_title,
page_description='',
page_keywords='',
alt_text=f"Screenshot of {page_uri}",
link_text='',
)
self.stats['screenshots_taken'] += 1
logger.debug(f"Screenshot captured: {page_uri} -> {md5_hash}")
except Exception as e:
logger.warning(f"Screenshot failed for {page_uri}: {e}")
def _get_extension(self, url: str, mime_type: str) -> str:
"""Determine file extension from URL or MIME type."""
from urllib.parse import urlparse
# Try from URL path
path = urlparse(url).path.lower()
for ext in ['.jpg', '.jpeg', '.png', '.gif', '.webp', '.svg', '.bmp',
'.mp4', '.webm', '.mov', '.avi', '.mkv']:
if path.endswith(ext):
return ext.lstrip('.')
# Try from MIME type
mime_map = {
'image/jpeg': 'jpg',
'image/png': 'png',
'image/gif': 'gif',
'image/webp': 'webp',
'image/svg+xml': 'svg',
'video/mp4': 'mp4',
'video/webm': 'webm',
}
for mt, ext in mime_map.items():
if mt in mime_type:
return ext
return 'bin'
async def main():
parser = argparse.ArgumentParser(
description="neopig - Neo Python Image Grabber",
epilog="Based on pig.py by Russell Ballestrini"
)
parser.add_argument(
"targets",
nargs="+",
help="Target URI(s) to crawl (e.g., https://example.com https://other.com)"
)
parser.add_argument(
"-k", "--keywords",
nargs="*",
default=[],
help="Keywords to tag media with (e.g., -k 'rick' 'morty')"
)
parser.add_argument(
"-m", "--mode",
choices=["text", "images", "videos", "media", "all"],
default="images",
help="Crawl mode (default: images)"
)
parser.add_argument(
"-d", "--depth",
type=int,
default=-1,
help="Crawl depth (-1 = unlimited, default: -1)"
)
parser.add_argument(
"-p", "--max-pages",
type=int,
default=-1,
help="Maximum pages to crawl (-1 = unlimited, default: -1)"
)
parser.add_argument(
"--db",
default="neopig.db",
help="Database path (default: neopig.db)"
)
parser.add_argument(
"--vault",
default="vault",
help="Vault storage path (default: vault)"
)
parser.add_argument(
"--no-download",
action="store_true",
help="Don't download media, just index URLs"
)
parser.add_argument(
"-v", "--verbose",
action="store_true",
help="Verbose output"
)
# Screenshot options (all off by default)
parser.add_argument(
"--screenshot",
action="store_true",
help="Enable page screenshots (requires uri2png)"
)
parser.add_argument(
"--screenshot-width",
type=int,
default=1280,
help="Screenshot viewport width in pixels (default: 1280)"
)
parser.add_argument(
"--screenshot-height",
type=int,
default=1024,
help="Screenshot viewport height in pixels (default: 1024)"
)
parser.add_argument(
"--screenshot-delay",
type=int,
default=1000,
help="Delay after page load in ms (default: 1000)"
)
args = parser.parse_args()
if args.verbose:
logging.getLogger().setLevel(logging.DEBUG)
# Map mode string to enum
mode_map = {
"text": CrawlMode.TEXT,
"images": CrawlMode.IMAGES,
"videos": CrawlMode.VIDEOS,
"media": CrawlMode.MEDIA,
"all": CrawlMode.ALL,
}
mode = mode_map[args.mode]
# Create screenshot config from args
screenshot_config = ScreenshotConfig(
enabled=args.screenshot,
width=args.screenshot_width,
height=args.screenshot_height,
delay=args.screenshot_delay,
)
if args.screenshot:
logger.info(f"Screenshots enabled: {screenshot_config.width}x{screenshot_config.height}, delay={screenshot_config.delay}ms")
# Initialize and run
pig = NeoPig(
db_path=args.db,
vault_path=args.vault,
screenshot_config=screenshot_config,
)
await pig.init()
# Load previously crawled media URIs to enable resume
crawled_media = await pig.db.get_crawled_media_uris()
if crawled_media:
logger.info(f"Resuming: {len(crawled_media)} media already crawled")
pig.seen_media = crawled_media
# Crawl all targets concurrently
async def crawl_target(target: str):
logger.info(f"=== Starting crawl: {target} ===")
return await pig.crawl(
target_uri=target,
keywords=args.keywords,
mode=mode,
depth=args.depth,
max_pages=args.max_pages,
download_media=not args.no_download,
)
await asyncio.gather(*[crawl_target(t) for t in args.targets])
if __name__ == "__main__":
asyncio.run(main())