Files
projects/Dockers/puppeteer-api/main.py
T
Bram 5fd11399b6
Build and Push Docker Images / build-and-push (push) Successful in 32s
cache tweaks
2025-05-01 13:14:08 +02:00

408 lines
13 KiB
Python

from fastapi import FastAPI, HTTPException, Header
from pyppeteer import launch
import os
import asyncio
import json
import re
import sqlite3
import time
from datetime import datetime, timedelta
from typing import Optional, Dict, List, Any
from urllib.parse import unquote
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.cron import CronTrigger
app = FastAPI()
# Get API key from environment variable
API_KEY = os.getenv('API_KEY')
if not API_KEY:
raise ValueError("API_KEY environment variable must be set")
# Get cache expiry time from environment variable (default: 36 hours)
CACHE_EXPIRY_HOURS = int(os.getenv('CACHE_EXPIRY_HOURS', '36'))
# Get cleanup cron schedule from environment variable (default: every day at 3 AM)
CLEANUP_CRON = os.getenv('CLEANUP_CRON', '0 3 * * *')
# Define custom user agent
CUSTOM_USER_AGENT = 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/69.0.3497.100 Safari/537.36'
# Initialize SQLite database
def init_db():
global DB_PATH
# Try to use the mounted volume first
db_path = '/db/cache.db'
db_dir = os.path.dirname(db_path)
# Check if directory exists and is writable
dir_writable = False
if os.path.exists(db_dir):
try:
test_file = os.path.join(db_dir, '.write_test')
with open(test_file, 'w') as f:
f.write('test')
os.remove(test_file)
dir_writable = True
except (IOError, PermissionError):
print(f"Directory {db_dir} exists but is not writable")
dir_writable = False
# If directory doesn't exist or isn't writable, try to create it
if not os.path.exists(db_dir) or not dir_writable:
try:
os.makedirs(db_dir, exist_ok=True)
# Test if we can write to the directory
test_file = os.path.join(db_dir, '.write_test')
with open(test_file, 'w') as f:
f.write('test')
os.remove(test_file)
print(f"Created directory: {db_dir}")
dir_writable = True
except Exception as e:
print(f"Warning: Could not create or write to directory {db_dir}: {e}")
# Fallback to using a local database file
db_path = 'cache.db'
print(f"Using local database file: {db_path}")
try:
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
cursor.execute('''
CREATE TABLE IF NOT EXISTS cache (
url TEXT,
route TEXT,
data TEXT,
timestamp INTEGER,
PRIMARY KEY (url, route)
)
''')
conn.commit()
conn.close()
print(f"Database initialized at {db_path}")
# Update the global DB_PATH
DB_PATH = db_path
except sqlite3.OperationalError as e:
print(f"Error initializing database at {db_path}: {e}")
# Fallback to using a local database file if the mounted volume has permission issues
db_path = 'cache.db'
print(f"Falling back to local database file: {db_path}")
try:
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
cursor.execute('''
CREATE TABLE IF NOT EXISTS cache (
url TEXT,
route TEXT,
data TEXT,
timestamp INTEGER,
PRIMARY KEY (url, route)
)
''')
conn.commit()
conn.close()
print(f"Local database initialized at {db_path}")
# Update the global DB_PATH
DB_PATH = db_path
except sqlite3.OperationalError as e2:
print(f"Error initializing local database: {e2}")
raise
# Define the database path
DB_PATH = '/db/cache.db'
# Get cached data if it exists and is not older than the expiry time
def get_cached_data(url, route):
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
cache_expiry = int(time.time()) - (CACHE_EXPIRY_HOURS * 60 * 60) # Convert hours to seconds
cursor.execute(
"SELECT data FROM cache WHERE url = ? AND route = ? AND timestamp > ?",
(url, route, cache_expiry)
)
result = cursor.fetchone()
conn.close()
if result:
print(f"Cache hit for {url} on route {route}")
return json.loads(result[0])
return None
# Save data to cache
def save_to_cache(url, route, data):
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
timestamp = int(time.time())
# Convert data to JSON string
data_json = json.dumps(data)
cursor.execute(
"INSERT OR REPLACE INTO cache (url, route, data, timestamp) VALUES (?, ?, ?, ?)",
(url, route, data_json, timestamp)
)
conn.commit()
conn.close()
print(f"Saved to cache: {url} on route {route}")
# Function to clean up old cache entries
def cleanup_old_cache_entries():
try:
print(f"Running scheduled cache cleanup (entries older than {CACHE_EXPIRY_HOURS} hours)")
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
# Calculate the timestamp for entries older than the expiry time
expiry_timestamp = int(time.time()) - (CACHE_EXPIRY_HOURS * 60 * 60)
# Get count of entries to be deleted
cursor.execute("SELECT COUNT(*) FROM cache WHERE timestamp < ?", (expiry_timestamp,))
count = cursor.fetchone()[0]
# Delete old entries
cursor.execute("DELETE FROM cache WHERE timestamp < ?", (expiry_timestamp,))
conn.commit()
conn.close()
print(f"Cache cleanup completed: {count} entries removed")
except Exception as e:
print(f"Error during cache cleanup: {e}")
# Initialize scheduler for periodic cache cleanup
scheduler = BackgroundScheduler()
scheduler.add_job(
cleanup_old_cache_entries,
CronTrigger.from_crontab(CLEANUP_CRON),
id='cache_cleanup_job',
replace_existing=True
)
# Initialize database on startup
init_db()
# Start the scheduler when the application starts
@app.on_event("startup")
def start_scheduler():
scheduler.start()
print(f"Cache cleanup scheduler started with cron: {CLEANUP_CRON}")
print(f"Cache entries will expire after {CACHE_EXPIRY_HOURS} hours")
# Shutdown the scheduler when the application stops
@app.on_event("shutdown")
def shutdown_scheduler():
scheduler.shutdown(wait=False)
print("Cache cleanup scheduler stopped")
async def wait_for_network_idle(page):
"""Wait until no network requests are in flight"""
await page.waitForNetworkIdle(idleTime=500, timeout=30000)
@app.head("/")
async def health_check():
return {"status": "ok"}
async def safe_browser_operation(url, operation_func):
"""Safely perform browser operations with proper cleanup"""
browser = None
try:
browser = await launch(
headless=True,
executablePath='/usr/bin/google-chrome',
args=['--no-sandbox', '--disable-setuid-sandbox'],
handleSIGINT=False,
handleSIGTERM=False,
handleSIGHUP=False
)
# Create new page with timeout
page = await browser.newPage()
page.setDefaultNavigationTimeout(30000)
# Set custom user agent
await page.setUserAgent(CUSTOM_USER_AGENT)
# Call the operation function that uses the page
result = await operation_func(page)
# Explicitly close the page
await page.close()
return result
finally:
# Ensure browser is closed properly
if browser:
try:
await browser.close()
except Exception as e:
print(f"Error closing browser: {e}")
# We don't re-raise here to avoid masking the original error
@app.get("/")
async def visit_url(url: str, x_api_key: Optional[str] = Header(None)):
# Validate API key
if not x_api_key or x_api_key != API_KEY:
raise HTTPException(status_code=401, detail="Invalid API key")
# Decode URL if it's encoded
decoded_url = unquote(url)
# Check cache first
cached_result = get_cached_data(decoded_url, "visit")
if cached_result:
return cached_result
try:
print(f"Visiting URL: {decoded_url}")
# Define the operation to perform with the browser
async def visit_operation(page):
try:
response = await page.goto(decoded_url, waitUntil='networkidle2', timeout=30000)
if not response:
print(f"Warning: No response object returned for {decoded_url}")
# Get page content
content = await page.content()
return {"status": "success", "content": content}
except Exception as e:
print(f"Error during page navigation: {e}")
# Try to get content anyway
try:
content = await page.content()
return {"status": "partial", "content": content, "error": str(e)}
except:
raise HTTPException(status_code=500, detail=f"Failed to get page content: {str(e)}")
# Perform the operation
result = await safe_browser_operation(decoded_url, visit_operation)
# Save to cache
save_to_cache(decoded_url, "visit", result)
return result
except Exception as e:
print(f"Error visiting URL {decoded_url}: {e}")
raise HTTPException(status_code=500, detail=str(e))
@app.get("/seo")
async def extract_seo(url: str, x_api_key: Optional[str] = Header(None)):
"""Extract SEO information from a website"""
# Validate API key
if not x_api_key or x_api_key != API_KEY:
raise HTTPException(status_code=401, detail="Invalid API key")
# Decode URL if it's encoded
decoded_url = unquote(url)
# Check cache first
cached_result = get_cached_data(decoded_url, "seo")
if cached_result:
return cached_result
try:
print(f"Extracting SEO from: {decoded_url}")
# Use the safe_browser_operation function to perform the browser operation
result = await safe_browser_operation(decoded_url, asyncio.run)
return result
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.get("/meta")
async def extract_meta_tags(url: str, x_api_key: Optional[str] = Header(None)):
"""Extract meta tags from a website"""
# Validate API key
if not x_api_key or x_api_key != API_KEY:
raise HTTPException(status_code=401, detail="Invalid API key")
# Decode URL if it's encoded
decoded_url = unquote(url)
# Check cache first
cached_result = get_cached_data(decoded_url, "meta")
if cached_result:
return cached_result
try:
print(f"Extracting meta tags from: {decoded_url}")
# Use the safe_browser_operation function to perform the browser operation
result = await safe_browser_operation(decoded_url, asyncio.run)
return result
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.get("/cache/clear")
async def clear_cache(x_api_key: Optional[str] = Header(None)):
"""Clear the entire cache database"""
# Validate API key
if not x_api_key or x_api_key != API_KEY:
raise HTTPException(status_code=401, detail="Invalid API key")
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
cursor.execute("DELETE FROM cache")
conn.commit()
conn.close()
return {"status": "success", "message": "Cache cleared successfully"}
@app.get("/cache/stats")
async def cache_stats(x_api_key: Optional[str] = Header(None)):
"""Get cache statistics"""
# Validate API key
if not x_api_key or x_api_key != API_KEY:
raise HTTPException(status_code=401, detail="Invalid API key")
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
# Get total entries
cursor.execute("SELECT COUNT(*) FROM cache")
total_entries = cursor.fetchone()[0]
# Get entries by route
cursor.execute("SELECT route, COUNT(*) FROM cache GROUP BY route")
routes = {route: count for route, count in cursor.fetchall()}
# Get recent entries (last 24 hours)
recent_timestamp = int(time.time()) - (24 * 60 * 60)
cursor.execute("SELECT COUNT(*) FROM cache WHERE timestamp > ?", (recent_timestamp,))
recent_entries = cursor.fetchone()[0]
# Get oldest entry timestamp
cursor.execute("SELECT MIN(timestamp) FROM cache")
oldest_timestamp = cursor.fetchone()[0]
oldest_date = datetime.fromtimestamp(oldest_timestamp).isoformat() if oldest_timestamp else None
# Get newest entry timestamp
cursor.execute("SELECT MAX(timestamp) FROM cache")
newest_timestamp = cursor.fetchone()[0]
newest_date = datetime.fromtimestamp(newest_timestamp).isoformat() if newest_timestamp else None
conn.close()
return {
"status": "success",
"stats": {
"total_entries": total_entries,
"entries_by_route": routes,
"recent_entries": recent_entries,
"oldest_entry": oldest_date,
"newest_entry": newest_date,
"cache_expiry_hours": CACHE_EXPIRY_HOURS,
"cleanup_schedule": CLEANUP_CRON
}
}
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)