776 lines
30 KiB
Python
776 lines
30 KiB
Python
"""
|
|
MTGJSON Data Manager Service
|
|
|
|
Handles downloading, unpacking, and upserting MTGJSON data into PostgreSQL.
|
|
Manages the complete data lifecycle from download to database upsert.
|
|
|
|
Features:
|
|
- Downloads required MTGJSON datasets from the API
|
|
- Unpacks gzip and zip files
|
|
- Converts JSON to PostgreSQL-compatible format
|
|
- Upserts data without overwriting existing entries
|
|
- Tracks refresh timestamps and status
|
|
- Provides health check data
|
|
|
|
Usage:
|
|
python -m app.services.mtgjson_manager [--force] [--refresh]
|
|
"""
|
|
|
|
import asyncio
|
|
import gzip
|
|
import json
|
|
import logging
|
|
import os
|
|
import re
|
|
import time
|
|
import zipfile
|
|
from datetime import datetime, timedelta
|
|
from pathlib import Path
|
|
from typing import Optional
|
|
|
|
import aiohttp
|
|
from sqlalchemy import text
|
|
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
|
|
from sqlalchemy.exc import IntegrityError
|
|
|
|
from app.core.database import mtg_async_session, mtg_engine
|
|
from app.core.settings import get_settings
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# MTGJSON API URLs
|
|
MTGJSON_BASE_URL = "https://mtgjson.com/api/v5"
|
|
REQUIRED_FILES = {
|
|
"AllPrintings.json.gz": MTGJSON_BASE_URL + "/AllPrintings.json",
|
|
"AllSetFiles.zip": MTGJSON_BASE_URL + "/AllSetFiles.zip",
|
|
"AllIdentifiers.json.gz": MTGJSON_BASE_URL + "/AllIdentifiers.json",
|
|
"CardTypes.json.gz": MTGJSON_BASE_URL + "/CardTypes.json",
|
|
"DeckList.json.gz": MTGJSON_BASE_URL + "/DeckList.json",
|
|
"Keywords.json.gz": MTGJSON_BASE_URL + "/Keywords.json",
|
|
"SetList.json.gz": MTGJSON_BASE_URL + "/SetList.json",
|
|
}
|
|
|
|
DATA_DIR = Path("/app/data/mtgjson")
|
|
REFRESH_LOG_TABLE = "mtg_refresh_log"
|
|
|
|
# Expected minimum file sizes (in bytes) for MTGJSON v5 files
|
|
EXPECTED_MIN_SIZES = {
|
|
"AllPrintings.json": 500 * 1024 * 1024, # 500 MB
|
|
"AllSetFiles.json": 10 * 1024 * 1024, # 10 MB
|
|
"AllIdentifiers.json": 100 * 1024 * 1024, # 100 MB
|
|
"CardTypes.json": 1 * 1024 * 1024, # 1 MB
|
|
"Keywords.json": 0.5 * 1024 * 1024, # 0.5 MB
|
|
"MagicSets.json": 50 * 1024 * 1024, # 50 MB
|
|
"MagicRoots.json": 1 * 1024 * 1024, # 1 MB
|
|
"SetTranslations.json": 5 * 1024 * 1024, # 5 MB
|
|
}
|
|
|
|
MAX_DOWNLOAD_RETRIES = 3
|
|
RETRY_DELAY_SECONDS = 60
|
|
|
|
|
|
class MTGJSONManager:
|
|
"""Manages MTGJSON data download, unpacking, and database upsert."""
|
|
|
|
def __init__(self, data_dir: Path = None):
|
|
"""Initialize MTGJSONManager.
|
|
|
|
Args:
|
|
data_dir: Path to data directory. If None, uses settings.DATA_DIR
|
|
"""
|
|
if data_dir is None:
|
|
settings = get_settings()
|
|
data_dir = settings.DATA_DIR
|
|
|
|
# Convert string to Path if needed
|
|
if isinstance(data_dir, str):
|
|
data_dir = Path(data_dir)
|
|
|
|
self.data_dir = data_dir
|
|
self.data_dir.mkdir(parents=True, exist_ok=True)
|
|
self.settings = get_settings()
|
|
|
|
async def validate_data_integrity(self) -> tuple[bool, list[str]]:
|
|
"""
|
|
Validate downloaded MTGJSON files for expected sizes.
|
|
|
|
Returns:
|
|
Tuple of (is_valid, list_of_issues)
|
|
"""
|
|
issues = []
|
|
|
|
for filename, min_size in EXPECTED_MIN_SIZES.items():
|
|
filepath = self.data_dir / filename
|
|
|
|
if not filepath.exists():
|
|
issues.append(f"Missing file: {filename}")
|
|
continue
|
|
|
|
actual_size = filepath.stat().st_size
|
|
|
|
if actual_size < min_size:
|
|
issues.append(
|
|
f"{filename}: {actual_size / (1024*1024):.1f} MB "
|
|
f"(minimum: {min_size / (1024*1024):.1f} MB)"
|
|
)
|
|
|
|
if issues:
|
|
logger.error(f"Data integrity check failed with {len(issues)} issues:")
|
|
for issue in issues:
|
|
logger.error(f" - {issue}")
|
|
return False, issues
|
|
|
|
logger.info("Data integrity check passed - all files meet minimum size requirements")
|
|
return True, []
|
|
|
|
async def cleanup_data_files(self) -> None:
|
|
"""Delete all downloaded MTGJSON data files."""
|
|
logger.info(f"Cleaning up data files in {self.data_dir}")
|
|
|
|
try:
|
|
for f in self.data_dir.glob("*"):
|
|
if f.is_file():
|
|
f.unlink()
|
|
logger.info(f"Deleted: {f.name}")
|
|
logger.info("Data cleanup complete")
|
|
except Exception as e:
|
|
logger.error(f"Failed to clean up data files: {e}")
|
|
raise
|
|
|
|
async def download_with_sanity_check(self) -> bool:
|
|
"""
|
|
Download MTGJSON files with sanity checking and retry logic.
|
|
|
|
Returns:
|
|
True if download and validation succeed, False otherwise
|
|
"""
|
|
logger.info("Starting MTGJSON download with sanity checks")
|
|
|
|
for attempt in range(1, MAX_DOWNLOAD_RETRIES + 1):
|
|
logger.info(f"Download attempt {attempt}/{MAX_DOWNLOAD_RETRIES}")
|
|
|
|
# Download files
|
|
download_success = await self.download_files()
|
|
|
|
if not download_success:
|
|
logger.error(f"Download failed on attempt {attempt}")
|
|
if attempt < MAX_DOWNLOAD_RETRIES:
|
|
logger.info(f"Waiting {RETRY_DELAY_SECONDS * attempt}s before retry...")
|
|
await asyncio.sleep(RETRY_DELAY_SECONDS * attempt)
|
|
await self.cleanup_data_files()
|
|
continue
|
|
return False
|
|
|
|
# Unpack files
|
|
unpack_success = await self.unpack_files()
|
|
|
|
if not unpack_success:
|
|
logger.error(f"Unpack failed on attempt {attempt}")
|
|
if attempt < MAX_DOWNLOAD_RETRIES:
|
|
logger.info(f"Waiting {RETRY_DELAY_SECONDS * attempt}s before retry...")
|
|
await asyncio.sleep(RETRY_DELAY_SECONDS * attempt)
|
|
await self.cleanup_data_files()
|
|
continue
|
|
return False
|
|
|
|
# Validate data integrity
|
|
is_valid, issues = await self.validate_data_integrity()
|
|
|
|
if is_valid:
|
|
logger.info(f"✓ Download and validation successful on attempt {attempt}")
|
|
return True
|
|
|
|
logger.warning(f"Data validation failed on attempt {attempt} with {len(issues)} issues")
|
|
|
|
if attempt < MAX_DOWNLOAD_RETRIES:
|
|
logger.info(f"Cleaning up and retrying in {RETRY_DELAY_SECONDS * attempt}s...")
|
|
await self.cleanup_data_files()
|
|
await asyncio.sleep(RETRY_DELAY_SECONDS * attempt)
|
|
else:
|
|
logger.error(f"Data validation failed after {MAX_DOWNLOAD_RETRIES} attempts")
|
|
return False
|
|
|
|
return False
|
|
|
|
async def download_files(self) -> bool:
|
|
"""Download all required MTGJSON files."""
|
|
logger.info(f"Starting MTGJSON data download to {self.data_dir}")
|
|
|
|
async with aiohttp.ClientSession() as session:
|
|
tasks = []
|
|
for filename, url in REQUIRED_FILES.items():
|
|
tasks.append(self._download_file(session, url, filename))
|
|
|
|
results = await asyncio.gather(*tasks, return_exceptions=True)
|
|
|
|
# Check if all downloads succeeded
|
|
success_count = sum(1 for r in results if r is True)
|
|
total = len(tasks)
|
|
|
|
if success_count == total:
|
|
logger.info(f"All {total} files downloaded successfully")
|
|
return True
|
|
else:
|
|
logger.error(f"Only {success_count}/{total} files downloaded")
|
|
return False
|
|
|
|
async def _download_file(self, session: aiohttp.ClientSession, url: str, filename: str) -> bool:
|
|
"""Download a single file."""
|
|
dest_path = self.data_dir / filename
|
|
|
|
# Skip if file already exists
|
|
if dest_path.exists() and dest_path.stat().st_size > 0:
|
|
logger.info(f"Skipping {filename} (already exists)")
|
|
return True
|
|
|
|
logger.info(f"Downloading {filename} from {url}")
|
|
|
|
try:
|
|
async with session.get(url) as response:
|
|
if response.status != 200:
|
|
logger.error(f"Failed to download {filename}: HTTP {response.status}")
|
|
return False
|
|
|
|
# Download with progress logging
|
|
content = b""
|
|
async for chunk in response.content.iter_chunked(8192):
|
|
content += chunk
|
|
|
|
# Write to file
|
|
dest_path.write_bytes(content)
|
|
logger.info(f"Downloaded {filename} ({len(content)} bytes)")
|
|
return True
|
|
|
|
except Exception as e:
|
|
logger.error(f"Download error for {filename}: {e}")
|
|
return False
|
|
|
|
async def unpack_files(self) -> bool:
|
|
"""Unpack all downloaded files."""
|
|
logger.info("Unpacking MTGJSON files")
|
|
|
|
unpacked_files = []
|
|
|
|
# Unpack gzip files
|
|
for gz_file in self.data_dir.glob("*.gz"):
|
|
if self._unpack_gzip(gz_file):
|
|
unpacked_files.append(gz_file.with_suffix(""))
|
|
|
|
# Unpack zip files
|
|
for zip_file in self.data_dir.glob("*.zip"):
|
|
if self._unpack_zip(zip_file):
|
|
unpacked_files.append(zip_file.with_suffix(""))
|
|
|
|
if len(unpacked_files) > 0:
|
|
logger.info(f"Unpacked {len(unpacked_files)} files")
|
|
return True
|
|
else:
|
|
logger.warning("No files were unpacked")
|
|
return False
|
|
|
|
def _unpack_gzip(self, gz_file: Path) -> bool:
|
|
"""Unpack a gzip file."""
|
|
dest_file = gz_file.with_suffix("")
|
|
|
|
logger.info(f"Unpacking {gz_file.name}")
|
|
|
|
try:
|
|
with gzip.open(gz_file, 'rt', encoding='utf-8') as f_in:
|
|
content = f_in.read()
|
|
|
|
dest_file.write_text(content, encoding='utf-8')
|
|
logger.info(f"Unpacked {gz_file.name} to {dest_file.name}")
|
|
return True
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to unpack {gz_file.name}: {e}")
|
|
return False
|
|
|
|
def _unpack_zip(self, zip_file: Path) -> bool:
|
|
"""Unpack a zip file."""
|
|
logger.info(f"Unpacking {zip_file.name}")
|
|
|
|
try:
|
|
with zipfile.ZipFile(zip_file, 'r') as zip_ref:
|
|
# Extract all files
|
|
zip_ref.extractall(self.data_dir)
|
|
|
|
# Remove the zip file after extraction
|
|
zip_file.unlink()
|
|
logger.info(f"Unpacked {zip_file.name}")
|
|
return True
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to unpack {zip_file.name}: {e}")
|
|
return False
|
|
|
|
async def upsert_data(self) -> dict[str, int]:
|
|
"""Upsert MTGJSON data into PostgreSQL.
|
|
|
|
Returns:
|
|
Dict with counts of sets and cards processed.
|
|
"""
|
|
logger.info("Upserting MTGJSON data into database")
|
|
|
|
counts = {
|
|
"sets": 0,
|
|
"cards": 0,
|
|
"identifiers": 0,
|
|
"card_types": 0,
|
|
"keywords": 0,
|
|
}
|
|
|
|
try:
|
|
# Process AllSetFiles directory
|
|
await self._upsert_sets()
|
|
counts["sets"] = await self._count_sets()
|
|
|
|
# Process AllPrintings.json
|
|
await self._upsert_cards()
|
|
counts["cards"] = await self._count_cards()
|
|
|
|
# Process AllIdentifiers.json
|
|
await self._upsert_identifiers()
|
|
counts["identifiers"] = await self._count_identifiers()
|
|
|
|
# Process CardTypes.json
|
|
await self._upsert_card_types()
|
|
counts["card_types"] = await self._count_card_types()
|
|
|
|
# Process Keywords.json
|
|
await self._upsert_keywords()
|
|
counts["keywords"] = await self._count_keywords()
|
|
|
|
logger.info(f"Data upsert complete: {counts}")
|
|
return counts
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to upsert data: {e}")
|
|
raise
|
|
|
|
async def _upsert_sets(self):
|
|
"""Upsert set data from AllSetFiles."""
|
|
allsetfiles_dir = self.data_dir / "allsetfiles"
|
|
|
|
if not allsetfiles_dir.exists():
|
|
logger.warning("AllSetFiles directory not found, skipping sets")
|
|
return
|
|
|
|
logger.info(f"Processing sets from {allsetfiles_dir}")
|
|
|
|
async with mtg_async_session() as session:
|
|
count = 0
|
|
|
|
for json_file in sorted(allsetfiles_dir.glob("*.json")):
|
|
try:
|
|
with open(json_file, 'r', encoding='utf-8') as f:
|
|
set_data = json.load(f)
|
|
|
|
if 'data' in set_data:
|
|
set_data = set_data['data']
|
|
|
|
# Get image URL
|
|
image_url = set_data.get('image', {}).get('png', set_data.get('image', {}).get('svg'))
|
|
|
|
# Upsert set
|
|
stmt = text("""
|
|
INSERT INTO mtg_sets (code, name, type, release_date, base_set_size,
|
|
total_size, is_foil_only, is_non_foil_only,
|
|
digital, icon_svg_url, parent_code, mtgo_code, image)
|
|
VALUES (:code, :name, :type, :release_date, :base_set_size,
|
|
:total_size, :is_foil_only, :is_non_foil_only,
|
|
:digital, :icon_svg_url, :parent_code, :mtgo_code, :image)
|
|
ON CONFLICT (code) DO UPDATE SET
|
|
name = EXCLUDED.name,
|
|
type = EXCLUDED.type,
|
|
release_date = EXCLUDED.release_date,
|
|
base_set_size = EXCLUDED.base_set_size,
|
|
total_size = EXCLUDED.total_size,
|
|
is_foil_only = EXCLUDED.is_foil_only,
|
|
is_non_foil_only = EXCLUDED.is_non_foil_only,
|
|
digital = EXCLUDED.digital,
|
|
icon_svg_url = EXCLUDED.icon_svg_url,
|
|
parent_code = EXCLUDED.parent_code,
|
|
mtgo_code = EXCLUDED.mtgo_code,
|
|
image = EXCLUDED.image,
|
|
updated_at = CURRENT_TIMESTAMP
|
|
""")
|
|
|
|
# Parse release date if present
|
|
release_date = None
|
|
release_date_str = set_data.get('releaseDate')
|
|
if release_date_str:
|
|
try:
|
|
release_date = datetime.fromisoformat(release_date_str.replace('Z', '+00:00'))
|
|
except (ValueError, AttributeError):
|
|
pass
|
|
|
|
await session.execute(stmt, {
|
|
'code': set_data.get('code'),
|
|
'name': set_data.get('name'),
|
|
'type': set_data.get('type'),
|
|
'release_date': release_date,
|
|
'base_set_size': set_data.get('baseSetSize'),
|
|
'total_size': set_data.get('totalSize'),
|
|
'is_foil_only': set_data.get('isFoilOnly'),
|
|
'is_non_foil_only': set_data.get('isNonFoilOnly'),
|
|
'digital': set_data.get('digital'),
|
|
'icon_svg_url': set_data.get('iconSvgUri'),
|
|
'parent_code': set_data.get('parentCode'),
|
|
'mtgo_code': set_data.get('mtgoCode'),
|
|
'image': image_url,
|
|
})
|
|
|
|
count += 1
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to upsert set {json_file.name}: {e}")
|
|
continue
|
|
|
|
await session.commit()
|
|
logger.info(f"Upserted {count} sets")
|
|
|
|
async def _upsert_cards(self):
|
|
"""Upsert card data from AllPrintings.json."""
|
|
allprintings_path = self.data_dir / "AllPrintings.json"
|
|
|
|
if not allprintings_path.exists():
|
|
logger.warning("AllPrintings.json not found, skipping cards")
|
|
return
|
|
|
|
logger.info("Processing cards from AllPrintings.json")
|
|
|
|
with open(allprintings_path, 'r', encoding='utf-8') as f:
|
|
data = json.load(f)
|
|
|
|
if 'data' not in data:
|
|
logger.error("Invalid AllPrintings.json structure: missing 'data' key")
|
|
return
|
|
|
|
mtg_data = data['data']
|
|
|
|
async with mtg_async_session() as session:
|
|
count = 0
|
|
|
|
for set_code, set_data in mtg_data.items():
|
|
if not isinstance(set_data, dict) or 'baseSetSize' not in set_data:
|
|
continue
|
|
|
|
cards = set_data.get('cards', [])
|
|
if not cards:
|
|
continue
|
|
|
|
# Get set_id
|
|
set_result = await session.execute(
|
|
text("SELECT id FROM mtg_sets WHERE code = :code"),
|
|
{'code': set_code}
|
|
)
|
|
set_row = set_result.fetchone()
|
|
|
|
if not set_row:
|
|
logger.warning(f"Set {set_code} not found in database")
|
|
continue
|
|
|
|
set_id = set_row[0]
|
|
|
|
# Upsert cards for this set
|
|
for card_data in cards:
|
|
# Ensure card has required fields
|
|
if 'name' not in card_data:
|
|
continue
|
|
|
|
try:
|
|
stmt = text("""
|
|
INSERT INTO mtg_cards (set_id, name, mana_cost, type_line, oracle_text,
|
|
power, toughness, rarity, layout, artist,
|
|
flavor_text, numbers, identifiers, images)
|
|
VALUES (:set_id, :name, :mana_cost, :type_line, :oracle_text,
|
|
:power, :toughness, :rarity, :layout, :artist,
|
|
:flavor_text, :numbers, :identifiers, :images)
|
|
ON CONFLICT (set_id, name) DO UPDATE SET
|
|
mana_cost = EXCLUDED.mana_cost,
|
|
type_line = EXCLUDED.type_line,
|
|
oracle_text = EXCLUDED.oracle_text,
|
|
power = EXCLUDED.power,
|
|
toughness = EXCLUDED.toughness,
|
|
rarity = EXCLUDED.rarity,
|
|
layout = EXCLUDED.layout,
|
|
artist = EXCLUDED.artist,
|
|
flavor_text = EXCLUDED.flavor_text,
|
|
numbers = EXCLUDED.numbers,
|
|
identifiers = EXCLUDED.identifiers,
|
|
images = EXCLUDED.images,
|
|
updated_at = CURRENT_TIMESTAMP
|
|
""")
|
|
|
|
# Serialize JSON fields
|
|
identifiers = json.dumps(card_data.get('identifiers', {}))
|
|
images = json.dumps(card_data.get('images', {}))
|
|
|
|
await session.execute(stmt, {
|
|
'set_id': set_id,
|
|
'name': card_data['name'],
|
|
'mana_cost': card_data.get('manaCost'),
|
|
'type_line': card_data.get('type'),
|
|
'oracle_text': card_data.get('text'),
|
|
'power': card_data.get('power'),
|
|
'toughness': card_data.get('toughness'),
|
|
'rarity': card_data.get('rarity'),
|
|
'layout': card_data.get('layout'),
|
|
'artist': card_data.get('artist'),
|
|
'flavor_text': card_data.get('flavorText'),
|
|
'numbers': str(card_data.get('numbers', '')),
|
|
'identifiers': identifiers,
|
|
'images': images,
|
|
})
|
|
|
|
count += 1
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to upsert card {card_data['name']} in set {set_code}: {e}")
|
|
continue
|
|
|
|
await session.commit()
|
|
logger.info(f"Upserted {count} cards")
|
|
|
|
async def _upsert_identifiers(self):
|
|
"""Upsert identifiers data."""
|
|
identifiers_path = self.data_dir / "AllIdentifiers.json"
|
|
|
|
if not identifiers_path.exists():
|
|
logger.warning("AllIdentifiers.json not found, skipping identifiers")
|
|
return
|
|
|
|
logger.info("Processing identifiers")
|
|
|
|
with open(identifiers_path, 'r', encoding='utf-8') as f:
|
|
identifiers = json.load(f)
|
|
|
|
# Store identifiers as JSON in a metadata table or as a file
|
|
identifiers_file = self.data_dir / "identifiers.json"
|
|
identifiers_file.write_text(json.dumps(identifiers, indent=2), encoding='utf-8')
|
|
logger.info(f"Saved identifiers to {identifiers_file}")
|
|
|
|
async def _upsert_card_types(self):
|
|
"""Upsert card types data."""
|
|
cardtypes_path = self.data_dir / "CardTypes.json"
|
|
|
|
if not cardtypes_path.exists():
|
|
logger.warning("CardTypes.json not found, skipping card types")
|
|
return
|
|
|
|
logger.info("Processing card types")
|
|
|
|
with open(cardtypes_path, 'r', encoding='utf-8') as f:
|
|
card_types = json.load(f)
|
|
|
|
# Store as JSON file
|
|
cardtypes_file = self.data_dir / "cardtypes.json"
|
|
cardtypes_file.write_text(json.dumps(card_types, indent=2), encoding='utf-8')
|
|
logger.info(f"Saved card types to {cardtypes_file}")
|
|
|
|
async def _upsert_keywords(self):
|
|
"""Upsert keywords data."""
|
|
keywords_path = self.data_dir / "Keywords.json"
|
|
|
|
if not keywords_path.exists():
|
|
logger.warning("Keywords.json not found, skipping keywords")
|
|
return
|
|
|
|
logger.info("Processing keywords")
|
|
|
|
with open(keywords_path, 'r', encoding='utf-8') as f:
|
|
keywords = json.load(f)
|
|
|
|
# Store as JSON file
|
|
keywords_file = self.data_dir / "keywords.json"
|
|
keywords_file.write_text(json.dumps(keywords, indent=2), encoding='utf-8')
|
|
logger.info(f"Saved keywords to {keywords_file}")
|
|
|
|
async def _count_sets(self) -> int:
|
|
"""Count total sets in database."""
|
|
async with mtg_async_session() as session:
|
|
result = await session.execute(text("SELECT COUNT(*) FROM mtg_sets"))
|
|
return result.scalar() or 0
|
|
|
|
async def _count_cards(self) -> int:
|
|
"""Count total cards in database."""
|
|
async with mtg_async_session() as session:
|
|
result = await session.execute(text("SELECT COUNT(*) FROM mtg_cards"))
|
|
return result.scalar() or 0
|
|
|
|
async def _count_identifiers(self) -> int:
|
|
"""Count identifiers file size."""
|
|
identifiers_file = self.data_dir / "identifiers.json"
|
|
if identifiers_file.exists():
|
|
return identifiers_file.stat().st_size
|
|
return 0
|
|
|
|
async def _count_card_types(self) -> int:
|
|
"""Count card types file size."""
|
|
cardtypes_file = self.data_dir / "cardtypes.json"
|
|
if cardtypes_file.exists():
|
|
return cardtypes_file.stat().st_size
|
|
return 0
|
|
|
|
async def _count_keywords(self) -> int:
|
|
"""Count keywords file size."""
|
|
keywords_file = self.data_dir / "keywords.json"
|
|
if keywords_file.exists():
|
|
return keywords_file.stat().st_size
|
|
return 0
|
|
|
|
async def log_refresh(self, status: str, counts: dict[str, int], duration: int, error: str = None):
|
|
"""Log refresh operation to database."""
|
|
async with mtg_async_session() as session:
|
|
stmt = text("""
|
|
INSERT INTO mtg_refresh_log (refresh_date, status, cards_updated, sets_updated,
|
|
identifiers_size, card_types_size, keywords_size,
|
|
error_message, duration_seconds)
|
|
VALUES (CURRENT_TIMESTAMP, :status, :cards, :sets,
|
|
:identifiers, :card_types, :keywords,
|
|
:error, :duration)
|
|
""")
|
|
|
|
await session.execute(stmt, {
|
|
'status': status,
|
|
'cards': counts.get('cards', 0),
|
|
'sets': counts.get('sets', 0),
|
|
'identifiers': counts.get('identifiers', 0),
|
|
'card_types': counts.get('card_types', 0),
|
|
'keywords': counts.get('keywords', 0),
|
|
'error': error,
|
|
'duration': duration,
|
|
})
|
|
|
|
await session.commit()
|
|
|
|
if status == "SUCCESS":
|
|
logger.info(f"Refresh logged: {counts}")
|
|
else:
|
|
logger.error(f"Refresh failed: {error}")
|
|
|
|
async def get_last_refresh(self) -> Optional[datetime]:
|
|
"""Get timestamp of last successful refresh."""
|
|
async with mtg_async_session() as session:
|
|
stmt = text("""
|
|
SELECT refresh_date FROM mtg_refresh_log
|
|
WHERE status = 'SUCCESS'
|
|
ORDER BY refresh_date DESC
|
|
LIMIT 1
|
|
""")
|
|
|
|
result = await session.execute(stmt)
|
|
row = result.fetchone()
|
|
|
|
return row[0] if row else None
|
|
|
|
def is_refresh_needed(self, last_refresh: Optional[datetime], interval_days: int = 7) -> bool:
|
|
"""Check if refresh is needed based on interval."""
|
|
if not last_refresh:
|
|
return True
|
|
|
|
days_since_refresh = (datetime.now() - last_refresh).days
|
|
return days_since_refresh >= interval_days
|
|
|
|
async def get_health_status(self) -> dict:
|
|
"""Get health status of MTGJSON data."""
|
|
try:
|
|
# Check if database is accessible
|
|
async with mtg_async_session() as session:
|
|
result = await session.execute(text("SELECT COUNT(*) FROM mtg_sets"))
|
|
sets_count = result.scalar() or 0
|
|
|
|
result = await session.execute(text("SELECT COUNT(*) FROM mtg_cards"))
|
|
cards_count = result.scalar() or 0
|
|
|
|
# Check if data files exist
|
|
data_files = {
|
|
"AllPrintings.json": self.data_dir / "AllPrintings.json",
|
|
"AllSetFiles": self.data_dir / "allsetfiles",
|
|
"AllIdentifiers.json": self.data_dir / "AllIdentifiers.json",
|
|
"CardTypes.json": self.data_dir / "CardTypes.json",
|
|
"Keywords.json": self.data_dir / "Keywords.json",
|
|
}
|
|
|
|
files_status = {}
|
|
for name, path in data_files.items():
|
|
if path.is_dir():
|
|
files_status[name] = "OK"
|
|
elif path.exists():
|
|
files_status[name] = "OK"
|
|
else:
|
|
files_status[name] = "MISSING"
|
|
|
|
# Get last refresh
|
|
last_refresh = await self.get_last_refresh()
|
|
|
|
return {
|
|
"status": "healthy" if sets_count > 0 and cards_count > 0 else "unhealthy",
|
|
"data": {
|
|
"sets_count": sets_count,
|
|
"cards_count": cards_count,
|
|
"last_refresh": last_refresh.isoformat() if last_refresh else None,
|
|
"files": files_status,
|
|
}
|
|
}
|
|
|
|
except Exception as e:
|
|
logger.error(f"Health check failed: {e}")
|
|
return {
|
|
"status": "unhealthy",
|
|
"error": str(e),
|
|
}
|
|
|
|
|
|
# Singleton instance
|
|
_manager_instance: Optional[MTGJSONManager] = None
|
|
|
|
|
|
def get_manager() -> MTGJSONManager:
|
|
"""Get or create MTGJSON manager instance."""
|
|
global _manager_instance
|
|
if _manager_instance is None:
|
|
_manager_instance = MTGJSONManager()
|
|
return _manager_instance
|
|
|
|
|
|
# CLI interface
|
|
if __name__ == "__main__":
|
|
import argparse
|
|
|
|
parser = argparse.ArgumentParser(description="MTGJSON Data Manager")
|
|
parser.add_argument("--force", action="store_true", help="Force download even if files exist")
|
|
parser.add_argument("--refresh", action="store_true", help="Run refresh cycle")
|
|
parser.add_argument("--health", action="store_true", help="Check health status")
|
|
parser.add_argument("--data-dir", type=Path, default=DATA_DIR, help="Data directory")
|
|
|
|
args = parser.parse_args()
|
|
|
|
async def main():
|
|
manager = MTGJSONManager(args.data_dir)
|
|
|
|
if args.health:
|
|
status = await manager.get_health_status()
|
|
print(json.dumps(status, indent=2))
|
|
elif args.refresh:
|
|
start_time = time.time()
|
|
try:
|
|
if args.force or not await manager.get_last_refresh():
|
|
if not await manager.download_with_sanity_check():
|
|
logger.error("Failed to download MTGJSON data after retries")
|
|
return
|
|
|
|
counts = await manager.upsert_data()
|
|
duration = int(time.time() - start_time)
|
|
await manager.log_refresh("SUCCESS", counts, duration)
|
|
|
|
except Exception as e:
|
|
duration = int(time.time() - start_time)
|
|
await manager.log_refresh("FAILED", {}, duration, str(e))
|
|
logger.error(f"Refresh failed: {e}")
|
|
|
|
else:
|
|
print("Usage: python -m app.services.mtgjson_manager [--refresh | --health]")
|
|
|
|
asyncio.run(main()) |