grant-radar / src /analyzer /crawler /scheduler.py
Riley
feat: Merge GPT-5 enhancements with working c72a240 base + restore crawler
cf9b3dc
Raw
History Blame Contribute Delete
10.8 kB
"""
Automated Grant Scraping Scheduler
Runs periodic crawls of Innovate UK and other funding sources to:
- Fetch new grants and update the database
- Mark grants as closed when deadlines pass
- Deduplicate entries
- Refresh the search index
Configuration:
- Run daily at 2 AM (configurable via CRAWLER_HOUR)
- Supports both APScheduler and system cron integration
"""
import logging
import os
from datetime import datetime, timedelta
from pathlib import Path
from typing import Optional, Dict, Any
from dataclasses import dataclass
logger = logging.getLogger(__name__)
# Configuration
DEFAULT_CRAWL_HOUR = int(os.getenv("CRAWLER_HOUR", "2")) # 2 AM
CRAWL_ENABLED = os.getenv("CRAWL_ENABLED", "false").lower() in ("true", "1", "yes")
SNAPSHOTS_DIR = Path(os.getenv("SNAPSHOTS_DIR", "data/snapshots"))
INDEX_PATH = Path(os.getenv("INDEX_PATH", "data/index/hybrid_index.pkl"))
@dataclass
class CrawlResult:
"""Result of a crawl operation."""
timestamp: datetime
new_grants: int
updated_grants: int
closed_grants: int
duplicates_removed: int
index_rebuilt: bool
error: Optional[str] = None
def to_dict(self) -> Dict[str, Any]:
return {
"timestamp": self.timestamp.isoformat(),
"new_grants": self.new_grants,
"updated_grants": self.updated_grants,
"closed_grants": self.closed_grants,
"duplicates_removed": self.duplicates_removed,
"index_rebuilt": self.index_rebuilt,
"error": self.error,
}
def deduplicate_grants(snapshots_dir: Path) -> int:
"""
Check for duplicate grants and keep only the most recent.
Returns:
Number of duplicates removed
"""
import json
if not snapshots_dir.exists():
return 0
# Group grants by ID (handles competition-XXXX variants)
grants_by_id: Dict[str, list] = {}
for json_file in snapshots_dir.glob("*.json"):
try:
with open(json_file, "r", encoding="utf-8") as f:
grant = json.load(f)
grant_id = grant.get("id") or json_file.stem
if grant_id not in grants_by_id:
grants_by_id[grant_id] = []
grants_by_id[grant_id].append({
"file": json_file,
"timestamp": json_file.stat().st_mtime,
"data": grant
})
except Exception as e:
logger.warning(f"Failed to read {json_file}: {e}")
continue
# Remove older duplicates
removed = 0
for grant_id, versions in grants_by_id.items():
if len(versions) > 1:
# Sort by timestamp, keep newest
versions.sort(key=lambda x: x["timestamp"], reverse=True)
for old_version in versions[1:]:
try:
old_version["file"].unlink()
removed += 1
logger.info(f"Removed duplicate: {old_version['file'].name}")
except Exception as e:
logger.warning(f"Failed to remove {old_version['file']}: {e}")
return removed
def mark_closed_grants(snapshots_dir: Path) -> int:
"""
Scan grants and mark those with passed deadlines as 'closed'.
Returns:
Number of grants marked as closed
"""
import json
from datetime import datetime
if not snapshots_dir.exists():
return 0
now = datetime.now()
closed_count = 0
for json_file in snapshots_dir.glob("*.json"):
try:
with open(json_file, "r", encoding="utf-8") as f:
grant = json.load(f)
# Check deadline
deadline_str = grant.get("close_date") or grant.get("deadline")
if not deadline_str:
continue
# Parse deadline
try:
# Handle ISO format with time
if "T" in deadline_str:
deadline = datetime.fromisoformat(deadline_str.replace("Z", "+00:00"))
else:
deadline = datetime.fromisoformat(deadline_str)
except ValueError:
continue
# If deadline passed and not marked closed, update it
if deadline < now and grant.get("status") != "closed":
grant["status"] = "closed"
grant["marked_closed_at"] = now.isoformat()
with open(json_file, "w", encoding="utf-8") as f:
json.dump(grant, f, indent=2)
logger.info(f"Marked as closed: {grant.get('title', json_file.stem)}")
closed_count += 1
except Exception as e:
logger.warning(f"Failed to process {json_file}: {e}")
continue
return closed_count
def rebuild_search_index(
snapshots_dir: Path = SNAPSHOTS_DIR,
output_path: Path = INDEX_PATH
) -> bool:
"""
Rebuild the search index from current data.
Returns:
True if successful
"""
try:
from ..search.hybrid_index import rebuild_index_from_data
from ..search.past_winners_integration import (
enrich_index_with_past_winners
)
from ..data_loader import load_past_winners
logger.info("Rebuilding search index...")
# Rebuild main index
idx = rebuild_index_from_data(
snapshots_dir=str(snapshots_dir),
output_path=str(output_path)
)
# Try to enrich with past winners
try:
past_winners = load_past_winners()
if past_winners:
enrich_index_with_past_winners(idx, past_winners)
logger.info(f"Enhanced index with {len(past_winners)} past winners")
except Exception as e:
logger.warning(f"Could not integrate past winners: {e}")
logger.info("Search index rebuilt successfully")
return True
except Exception as e:
logger.error(f"Failed to rebuild search index: {e}")
return False
def run_crawl_cycle() -> CrawlResult:
"""
Run a complete crawl and maintenance cycle.
This is the main entry point for scheduled crawls.
"""
import asyncio
logger.info("Starting grant crawler cycle...")
start_time = datetime.now()
try:
# Step 1: Discover and fetch new grants from Innovate UK
new_grants = 0
try:
from ...crawler.discover_grants import discover_and_fetch_grants
logger.info("Discovering new grants from Innovate UK...")
total_discovered, newly_fetched, new_files = asyncio.run(
discover_and_fetch_grants(SNAPSHOTS_DIR, skip_existing=True)
)
new_grants = newly_fetched
logger.info(
f"Grant discovery: {total_discovered} total, "
f"{newly_fetched} newly fetched"
)
if new_files:
logger.info(f"New grants: {', '.join(new_files)}")
except Exception as e:
logger.error(f"Grant discovery failed: {e}", exc_info=True)
# Continue with other steps even if discovery fails
# Step 2: Check for duplicates
duplicates = deduplicate_grants(SNAPSHOTS_DIR)
# Step 3: Mark closed grants
closed = mark_closed_grants(SNAPSHOTS_DIR)
# Step 4: Rebuild index
index_rebuilt = rebuild_search_index()
result = CrawlResult(
timestamp=start_time,
new_grants=new_grants,
updated_grants=0,
closed_grants=closed,
duplicates_removed=duplicates,
index_rebuilt=index_rebuilt,
)
logger.info(f"Crawl cycle complete: {result}")
return result
except Exception as e:
logger.error(f"Crawl cycle failed: {e}", exc_info=True)
return CrawlResult(
timestamp=start_time,
new_grants=0,
updated_grants=0,
closed_grants=0,
duplicates_removed=0,
index_rebuilt=False,
error=str(e)
)
def setup_scheduler():
"""
Set up APScheduler for daily crawls.
Usage:
from analyzer.crawler.scheduler import setup_scheduler
scheduler = setup_scheduler()
scheduler.start()
"""
try:
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.cron import CronTrigger
if not CRAWL_ENABLED:
logger.info("Crawl scheduler disabled (set CRAWL_ENABLED=true to enable)")
return None
scheduler = BackgroundScheduler()
# Schedule daily crawl at DEFAULT_CRAWL_HOUR (default 2 AM)
trigger = CronTrigger(hour=DEFAULT_CRAWL_HOUR, minute=0)
scheduler.add_job(
run_crawl_cycle,
trigger=trigger,
id="grant_crawler",
name="Daily grant crawl and index refresh",
replace_existing=True
)
logger.info(
f"Scheduler configured: Daily crawl at {DEFAULT_CRAWL_HOUR}:00 "
f"(set CRAWL_ENABLED=true to start)"
)
return scheduler
except ImportError:
logger.warning(
"APScheduler not installed. "
"Install with: pip install apscheduler"
)
return None
def register_cron_job():
"""
Register a system cron job for daily crawls.
Useful as alternative to APScheduler for production deployments.
Example cron line (runs daily at 2 AM):
0 2 * * * cd /path/to/grant-analyst && python -m analyzer.crawler.scheduler
"""
import subprocess
import platform
if platform.system() == "Windows":
logger.warning("Cron registration only supported on Unix-like systems")
return False
try:
script_path = Path(__file__).parent.parent.parent / "crawler" / "scheduler.py"
cron_line = f"0 {DEFAULT_CRAWL_HOUR} * * * python {script_path}"
# This is a guide - actual registration depends on system setup
logger.info(f"Add this to crontab for daily crawls:\n{cron_line}")
return True
except Exception as e:
logger.error(f"Failed to register cron job: {e}")
return False
# ============================================================================
# Entry Points
# ============================================================================
if __name__ == "__main__":
import sys
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s - %(name)s - %(levelname)s - %(message)s"
)
# Run single crawl cycle
result = run_crawl_cycle()
print(f"\nCrawl Result: {result.to_dict()}")
if result.error:
sys.exit(1)
sys.exit(0)