""" MTGJSON Data Manager - Local File Processing Reads JSON/PSQL files from a mounted volume and upserts them into PostgreSQL. Handles AllDeckFiles.zip extraction and processing. Includes download functionality for container init and manual refresh. """ import json import logging import zipfile import tempfile import shutil from pathlib import Path from datetime import datetime from typing import Optional, Dict, List, Tuple import re import aiohttp import asyncio import gzip from sqlalchemy import text from app.core.database import mtg_async_session logger = logging.getLogger(__name__) class MTGJSONManager: """Process local MTGJSON files and upsert into PostgreSQL.""" def __init__(self, data_dir: Path = None): if data_dir is None: data_dir = Path("/app/data") if isinstance(data_dir, str): data_dir = Path(data_dir) self.data_dir = data_dir self.data_dir.mkdir(parents=True, exist_ok=True) async def get_last_refresh(self) -> Optional[datetime]: """Get timestamp of last successful refresh.""" async with mtg_async_session() as session: result = await session.execute(text(""" SELECT refresh_date FROM mtg_refresh_log WHERE status = 'SUCCESS' ORDER BY refresh_date DESC LIMIT 1 """)) row = result.fetchone() return row[0] if row else None async def get_health_status(self) -> dict: """Get health status.""" try: 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 files_status = {} for f in self.data_dir.iterdir(): if f.is_file(): if f.suffix == '.json': files_status[f.stem + '.json'] = 'OK' elif f.suffix == '.psql': files_status[f.stem + '.psql'] = 'OK' elif f.suffix == '.zip': files_status[f.stem + '.zip'] = 'OK' elif f.is_dir(): files_status[f.name] = 'OK' last_refresh = await self.get_last_refresh() required = await self.check_required_files() # App is healthy if: # 1. Cards and sets are loaded # 2. All required files are present is_healthy = ( sets_count > 0 and cards_count > 0 and required['all_present'] ) return { "status": "healthy" if is_healthy else "unhealthy", "data": { "sets_count": sets_count, "cards_count": cards_count, "last_refresh": last_refresh.isoformat() if last_refresh else None, "files": files_status, "required_files": required, } } except Exception as e: return {"status": "unhealthy", "error": str(e)} async def check_required_files(self) -> dict: """Check if all required files are present.""" available = set() for item in self.data_dir.iterdir(): if item.is_file(): available.add(item.name) elif item.is_dir(): available.add(item.name) required_files = { 'AllPrintings.json': 'AllPrintings.json', 'AllPrintings.psql': 'AllPrintings.psql', 'AllIdentifiers.json': 'AllIdentifiers.json', 'Keywords.json': 'Keywords.json', 'CardTypes.json': 'CardTypes.json', 'AllDeckFiles.zip': 'AllDeckFiles.zip', } result = { 'all_present': True, 'missing': [], 'found': [], 'details': {} } # Check AllPrintings (either .json or .psql) if 'AllPrintings.json' in available: result['details']['AllPrintings'] = 'OK (.json)' result['found'].append('AllPrintings.json') elif 'AllPrintings.psql' in available: result['details']['AllPrintings'] = 'OK (.psql)' result['found'].append('AllPrintings.psql') else: result['details']['AllPrintings'] = 'MISSING' result['missing'].append('AllPrintings.json or AllPrintings.psql') result['all_present'] = False # Check AllIdentifiers.json if 'AllIdentifiers.json' in available: result['details']['AllIdentifiers'] = 'OK' result['found'].append('AllIdentifiers.json') else: result['details']['AllIdentifiers'] = 'MISSING' result['missing'].append('AllIdentifiers.json') result['all_present'] = False # Check Keywords.json if 'Keywords.json' in available: result['details']['Keywords'] = 'OK' result['found'].append('Keywords.json') else: result['details']['Keywords'] = 'MISSING' result['missing'].append('Keywords.json') result['all_present'] = False # Check CardTypes.json if 'CardTypes.json' in available: result['details']['CardTypes'] = 'OK' result['found'].append('CardTypes.json') else: result['details']['CardTypes'] = 'MISSING' result['missing'].append('CardTypes.json') result['all_present'] = False # Check AllDeckFiles.zip if 'AllDeckFiles.zip' in available: result['details']['AllDeckFiles'] = 'OK' result['found'].append('AllDeckFiles.zip') else: result['details']['AllDeckFiles'] = 'MISSING' result['missing'].append('AllDeckFiles.zip') result['all_present'] = False return result async def upsert_all(self) -> dict: """Upsert all data from local files.""" counts = {} # Check what files exist available = set() for item in self.data_dir.iterdir(): if item.is_file(): available.add(item.name) elif item.is_dir(): available.add(item.name) logger.info(f"Found files: {available}") # Upsert each file type if "AllSetFiles" in available: counts["sets"] = await self._upsert_sets() # Handle AllPrintings - can be .json or .psql if "AllPrintings.json" in available: counts["cards"] = await self._upsert_cards() elif "AllPrintings.psql" in available: counts["cards"] = await self._upsert_cards_from_psql() if "AllIdentifiers.json" in available: counts["identifiers"] = await self._upsert_identifiers() if "CardTypes.json" in available: counts["card_types"] = await self._upsert_card_types() if "Keywords.json" in available: counts["keywords"] = await self._upsert_keywords() if "SetList.json" in available: counts["set_list"] = await self._upsert_set_list() # Handle deck files - can be .zip or .json if "AllDeckFiles.zip" in available: counts["deck_list"] = await self._upsert_deck_files_from_zip() elif "DeckList.json" in available: counts["deck_list"] = await self._upsert_deck_list() return counts async def _upsert_sets(self) -> int: """Upsert sets from AllSetFiles directory.""" allsetfiles = self.data_dir / "AllSetFiles" if not allsetfiles.exists(): logger.warning("AllSetFiles directory not found") return 0 logger.info(f"Processing sets from {allsetfiles}") count = 0 async with mtg_async_session() as session: for json_file in sorted(allsetfiles.glob("*.json")): logger.info(f"Processing {json_file.name}") try: with open(json_file, 'r', encoding='utf-8') as f: data = json.load(f) if not isinstance(data, dict): logger.warning(f"{json_file.name} is not a dict") continue # Upsert sets sets = data.get("sets", []) for set_data in sets: if not isinstance(set_data, dict): continue await session.execute(text(""" INSERT INTO mtg_sets (set_name, code, name, type, release_date, scryfall_uri, card_count, tokens_count, is_foil_only, is_non_foil_only) VALUES ( :set_name, :code, :name, :type, :release_date, :scryfall_uri, :card_count, :tokens_count, :is_foil_only, :is_non_foil_only ) ON CONFLICT (set_name) DO UPDATE SET code = EXCLUDED.code, name = EXCLUDED.name, type = EXCLUDED.type, release_date = EXCLUDED.release_date, scryfall_uri = EXCLUDED.scryfall_uri, card_count = EXCLUDED.card_count, tokens_count = EXCLUDED.tokens_count, is_foil_only = EXCLUDED.is_foil_only, is_non_foil_only = EXCLUDED.is_non_foil_only """), { 'set_name': set_data.get('name'), 'code': set_data.get('code'), 'name': set_data.get('name'), 'type': set_data.get('type'), 'release_date': set_data.get('release_date'), 'scryfall_uri': set_data.get('scryfall_uri'), 'card_count': set_data.get('card_count'), 'tokens_count': set_data.get('tokens_count'), 'is_foil_only': set_data.get('is_foil_only'), 'is_non_foil_only': set_data.get('is_non_foil_only'), }) count += 1 await session.commit() logger.info(f"Upserted {count} sets from {json_file.name}") except Exception as e: logger.error(f"Failed to process {json_file.name}: {e}") await session.rollback() continue return count async def _upsert_cards(self) -> int: """Upsert cards from AllPrintings.json.""" path = self.data_dir / "AllPrintings.json" if not path.exists(): logger.warning("AllPrintings.json not found") return 0 logger.info(f"Processing cards from {path}") count = 0 with open(path, 'r', encoding='utf-8') as f: data = json.load(f) cards = data.get("cards", []) async with mtg_async_session() as session: for card in cards: if not isinstance(card, dict): continue try: await session.execute(text(""" INSERT INTO mtg_cards ( card_name, mtgo_id, card_type, set_name, rarity, artist, number, language, mana_cost, text, power, toughness, loyalty, colors, color_identity, produced_mana, legalities, original_type, foreign_data, rulings, hand_modifier, life_modifier, side_names ) VALUES ( :card_name, :mtgo_id, :card_type, :set_name, :rarity, :artist, :number, :language, :mana_cost, :text, :power, :toughness, :loyalty, :colors, :color_identity, :produced_mana, :legalities, :original_type, :foreign_data, :rulings, :hand_modifier, :life_modifier, :side_names ) ON CONFLICT (card_name, set_name, mtgo_id) DO UPDATE SET card_type = EXCLUDED.card_type, rarity = EXCLUDED.rarity, artist = EXCLUDED.artist, number = EXCLUDED.number, language = EXCLUDED.language, mana_cost = EXCLUDED.mana_cost, text = EXCLUDED.text, power = EXCLUDED.power, toughness = EXCLUDED.toughness, loyalty = EXCLUDED.loyalty, colors = EXCLUDED.colors, color_identity = EXCLUDED.color_identity, produced_mana = EXCLUDED.produced_mana, legalities = EXCLUDED.legalities, original_type = EXCLUDED.original_type, foreign_data = EXCLUDED.foreign_data, rulings = EXCLUDED.rulings, hand_modifier = EXCLUDED.hand_modifier, life_modifier = EXCLUDED.life_modifier, side_names = EXCLUDED.side_names """), { 'card_name': card.get('name'), 'mtgo_id': card.get('mtgoId'), 'card_type': card.get('type'), 'set_name': card.get('setName'), 'rarity': card.get('rarity'), 'artist': card.get('artist'), 'number': card.get('number'), 'language': card.get('language'), 'mana_cost': json.dumps(card.get('manaCost')), 'text': card.get('text'), 'power': card.get('power'), 'toughness': card.get('toughness'), 'loyalty': card.get('loyalty'), 'colors': json.dumps(card.get('colors')), 'color_identity': json.dumps(card.get('colorIdentity')), 'produced_mana': json.dumps(card.get('producedMana')), 'legalities': json.dumps(card.get('legalities')), 'original_type': card.get('originalType'), 'foreign_data': json.dumps(card.get('foreignData')), 'rulings': json.dumps(card.get('rulings')), 'hand_modifier': card.get('handModifier'), 'life_modifier': card.get('lifeModifier'), 'side_names': json.dumps(card.get('sideNames')), }) count += 1 except Exception as e: logger.error(f"Failed to upsert card: {e}") await session.rollback() continue logger.info(f"Upserted {count} cards") return count async def _upsert_cards_from_psql(self) -> int: """Upsert cards from AllPrintings.psql file. PSQL files contain SQL INSERT statements. We parse them to extract card data. """ path = self.data_dir / "AllPrintings.psql" if not path.exists(): logger.warning("AllPrintings.psql not found") return 0 logger.info(f"Processing cards from PSQL file: {path}") cards = self._parse_psql_file(path) async with mtg_async_session() as session: count = 0 for card in cards: if not isinstance(card, dict): continue try: await session.execute(text(""" INSERT INTO mtg_cards ( card_name, mtgo_id, card_type, set_name, rarity, artist, number, language, mana_cost, text, power, toughness, loyalty, colors, color_identity, produced_mana, legalities, original_type, foreign_data, rulings, hand_modifier, life_modifier, side_names ) VALUES ( :card_name, :mtgo_id, :card_type, :set_name, :rarity, :artist, :number, :language, :mana_cost, :text, :power, :toughness, :loyalty, :colors, :color_identity, :produced_mana, :legalities, :original_type, :foreign_data, :rulings, :hand_modifier, :life_modifier, :side_names ) ON CONFLICT (card_name, set_name, mtgo_id) DO UPDATE SET card_type = EXCLUDED.card_type, rarity = EXCLUDED.rarity, artist = EXCLUDED.artist, number = EXCLUDED.number, language = EXCLUDED.language, mana_cost = EXCLUDED.mana_cost, text = EXCLUDED.text, power = EXCLUDED.power, toughness = EXCLUDED.toughness, loyalty = EXCLUDED.loyalty, colors = EXCLUDED.colors, color_identity = EXCLUDED.color_identity, produced_mana = EXCLUDED.produced_mana, legalities = EXCLUDED.legalities, original_type = EXCLUDED.original_type, foreign_data = EXCLUDED.foreign_data, rulings = EXCLUDED.rulings, hand_modifier = EXCLUDED.hand_modifier, life_modifier = EXCLUDED.life_modifier, side_names = EXCLUDED.side_names """), { 'card_name': card.get('name'), 'mtgo_id': card.get('mtgoId'), 'card_type': card.get('type'), 'set_name': card.get('setName'), 'rarity': card.get('rarity'), 'artist': card.get('artist'), 'number': card.get('number'), 'language': card.get('language'), 'mana_cost': json.dumps(card.get('manaCost')), 'text': card.get('text'), 'power': card.get('power'), 'toughness': card.get('toughness'), 'loyalty': card.get('loyalty'), 'colors': json.dumps(card.get('colors')), 'color_identity': json.dumps(card.get('colorIdentity')), 'produced_mana': json.dumps(card.get('producedMana')), 'legalities': json.dumps(card.get('legalities')), 'original_type': card.get('originalType'), 'foreign_data': json.dumps(card.get('foreignData')), 'rulings': json.dumps(card.get('rulings')), 'hand_modifier': card.get('handModifier'), 'life_modifier': card.get('lifeModifier'), 'side_names': json.dumps(card.get('sideNames')), }) count += 1 except Exception as e: logger.error(f"Failed to upsert card: {e}") continue await session.commit() logger.info(f"Upserted {count} cards from PSQL") return count def _parse_psql_file(self, filepath: Path) -> list: """Parse a PSQL file containing SQL INSERT statements to extract card data.""" cards = [] try: with open(filepath, 'r', encoding='utf-8') as f: content = f.read() # Find all INSERT INTO mtg_cards statements insert_pattern = r"INSERT INTO mtg_cards\s*\(([^)]+)\)\s*VALUES\s*\(([^;]+)\);" matches = re.findall(insert_pattern, content, re.IGNORECASE | re.DOTALL) for columns_str, values_str in matches: # Parse column names columns = [col.strip().strip("'\"") for col in columns_str.split(',')] # Parse values (simplified parser) values = self._parse_sql_values(values_str) if len(columns) == len(values): card = dict(zip(columns, values)) # Try to parse JSON fields json_fields = ['manaCost', 'colors', 'colorIdentity', 'producedMana', 'legalities', 'foreignData', 'rulings', 'sideNames'] for field in json_fields: if field in card and isinstance(card[field], str): try: card[field] = json.loads(card[field]) except: pass cards.append(card) except Exception as e: logger.error(f"Error parsing PSQL file {filepath}: {e}") return cards def _parse_sql_values(self, values_str: str) -> list: """Parse SQL VALUES clause to extract individual values.""" values = [] current_value = "" in_single_quote = False in_double_quote = False escape_next = False for char in values_str: if escape_next: current_value += char escape_next = False continue if char == '\\' and not in_single_quote: escape_next = True continue if char == "'" and not in_double_quote: in_single_quote = not in_single_quote current_value += char continue if char == '"' and not in_single_quote: in_double_quote = not in_double_quote current_value += char continue if char == ',' and not in_single_quote and not in_double_quote: values.append(self._clean_sql_value(current_value.strip())) current_value = "" continue current_value += char # Add last value if current_value.strip(): values.append(self._clean_sql_value(current_value.strip())) return values def _clean_sql_value(self, value: str) -> str: """Clean and parse a SQL value.""" value = value.strip() # Remove quotes if (value.startswith("'") and value.endswith("'")) or \ (value.startswith('"') and value.endswith('"')): value = value[1:-1] # Check for NULL if value.upper() == "NULL": return None # Try to parse as JSON if value.startswith("{") or value.startswith("["): try: return json.loads(value.replace("NULL", "null")) except: pass return value async def _upsert_identifiers(self) -> int: """Upsert identifiers from AllIdentifiers.json.""" path = self.data_dir / "AllIdentifiers.json" if not path.exists(): logger.warning("AllIdentifiers.json not found") return 0 logger.info(f"Processing identifiers from {path}") with open(path, 'r', encoding='utf-8') as f: data = json.load(f) identifiers = data.get("identifiers", {}) async with mtg_async_session() as session: count = 0 for identifier_id, identifier in identifiers.items(): if not isinstance(identifier, dict): continue try: await session.execute(text(""" INSERT INTO mtg_identifiers (identifier_id, identifier_data, created_at) VALUES (:identifier_id, :identifier_data, CURRENT_TIMESTAMP) ON CONFLICT (identifier_id) DO UPDATE SET identifier_data = EXCLUDED.identifier_data, updated_at = CURRENT_TIMESTAMP """), { 'identifier_id': identifier_id, 'identifier_data': json.dumps(identifier), }) count += 1 except Exception as e: logger.error(f"Failed to upsert identifier {identifier_id}: {e}") continue await session.commit() logger.info(f"Upserted {count} identifiers") return count async def _upsert_card_types(self) -> int: """Upsert card types from CardTypes.json.""" path = self.data_dir / "CardTypes.json" if not path.exists(): logger.warning("CardTypes.json not found") return 0 logger.info("Processing card types") with open(path, 'r', encoding='utf-8') as f: card_types = json.load(f) async with mtg_async_session() as session: count = 0 for card_type in card_types: if not isinstance(card_type, dict): continue try: await session.execute(text(""" INSERT INTO mtg_card_types (card_type) VALUES (:card_type) ON CONFLICT (card_type) DO NOTHING """), { 'card_type': card_type.get('cardType'), }) count += 1 except Exception as e: logger.error(f"Failed to upsert card type: {e}") continue await session.commit() logger.info(f"Upserted {count} card types") return count async def _upsert_keywords(self) -> int: """Upsert keywords from Keywords.json.""" path = self.data_dir / "Keywords.json" if not path.exists(): logger.warning("Keywords.json not found") return 0 logger.info("Processing keywords") with open(path, 'r', encoding='utf-8') as f: keywords = json.load(f) async with mtg_async_session() as session: count = 0 for keyword in keywords: if not isinstance(keyword, dict): continue try: await session.execute(text(""" INSERT INTO mtg_keywords (keyword) VALUES (:keyword) ON CONFLICT (keyword) DO NOTHING """), { 'keyword': keyword.get('keyword'), }) count += 1 except Exception as e: logger.error(f"Failed to upsert keyword: {e}") continue await session.commit() logger.info(f"Upserted {count} keywords") return count async def _upsert_set_list(self) -> int: """Upsert set list from SetList.json.""" path = self.data_dir / "SetList.json" if not path.exists(): logger.warning("SetList.json not found") return 0 logger.info("Processing set list") with open(path, 'r', encoding='utf-8') as f: set_list = json.load(f) async with mtg_async_session() as session: count = 0 for item in set_list: if not isinstance(item, dict): continue try: await session.execute(text(""" INSERT INTO mtg_set_list (set_name, code, release_date, scryfall_uri) VALUES (:set_name, :code, :release_date, :scryfall_uri) ON CONFLICT (set_name) DO UPDATE SET code = EXCLUDED.code, release_date = EXCLUDED.release_date, scryfall_uri = EXCLUDED.scryfall_uri """), { 'set_name': item.get('name'), 'code': item.get('code'), 'release_date': item.get('release_date'), 'scryfall_uri': item.get('scryfall_uri'), }) count += 1 except Exception as e: logger.error(f"Failed to upsert set list item: {e}") continue await session.commit() logger.info(f"Upserted {count} set list items") return count async def _upsert_deck_list(self) -> int: """Upsert deck list from DeckList.json.""" path = self.data_dir / "DeckList.json" if not path.exists(): logger.warning("DeckList.json not found") return 0 logger.info("Processing deck list") with open(path, 'r', encoding='utf-8') as f: deck_list = json.load(f) async with mtg_async_session() as session: count = 0 for item in deck_list: if not isinstance(item, dict): continue try: await session.execute(text(""" INSERT INTO mtg_deck_list (list_id, list_name, list_year, list_date, list_format) VALUES (:list_id, :list_name, :list_year, :list_date, :list_format) ON CONFLICT (list_id) DO UPDATE SET list_name = EXCLUDED.list_name, list_year = EXCLUDED.list_year, list_date = EXCLUDED.list_date, list_format = EXCLUDED.list_format """), { 'list_id': item.get('listId'), 'list_name': item.get('listName'), 'list_year': item.get('listYear'), 'list_date': item.get('listDate'), 'list_format': item.get('listFormat'), }) count += 1 except Exception as e: logger.error(f"Failed to upsert deck list item: {e}") continue await session.commit() logger.info(f"Upserted {count} deck list items") return count async def _upsert_deck_files_from_zip(self) -> int: """Unzip and upsert deck files from AllDeckFiles.zip.""" zip_path = self.data_dir / "AllDeckFiles.zip" if not zip_path.exists(): logger.warning("AllDeckFiles.zip not found") return 0 logger.info(f"Processing deck files from zip: {zip_path}") json_files = self._extract_zip_files(zip_path) async with mtg_async_session() as session: count = 0 for json_file in json_files: logger.info(f"Processing deck file: {json_file.name}") try: with open(json_file, 'r', encoding='utf-8') as f: deck_data = json.load(f) if not isinstance(deck_data, dict): logger.warning(f"{json_file.name} is not a dict") continue # Upsert deck await session.execute(text(""" INSERT INTO mtg_deck_list (list_id, list_name, list_year, list_date, list_format) VALUES ( :list_id, :list_name, :list_year, :list_date, :list_format ) ON CONFLICT (list_id) DO UPDATE SET list_name = EXCLUDED.list_name, list_year = EXCLUDED.list_year, list_date = EXCLUDED.list_date, list_format = EXCLUDED.list_format """), { 'list_id': deck_data.get('listId', str(hash(str(deck_data)))), 'list_name': deck_data.get('name', json_file.stem), 'list_year': deck_data.get('year', 'unknown'), 'list_date': deck_data.get('date', 'unknown'), 'list_format': deck_data.get('format', 'unknown'), }) count += 1 except Exception as e: logger.error(f"Failed to upsert deck file {json_file.name}: {e}") continue await session.commit() logger.info(f"Upserted {count} deck files from zip") return count def _extract_zip_files(self, zip_path: Path) -> list: """Extract JSON files from a ZIP archive.""" json_files = [] try: with zipfile.ZipFile(zip_path, 'r') as zip_ref: # Create temporary directory for extraction with tempfile.TemporaryDirectory() as temp_dir: zip_ref.extractall(temp_dir) # Find all JSON files recursively for file_path in Path(temp_dir).rglob("*.json"): json_files.append(file_path) except Exception as e: logger.error(f"Error extracting zip file {zip_path}: {e}") return json_files async def log_refresh(self, status: str, counts: dict[str, int], duration: int, error: str = None): """Log refresh operation.""" 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, sets_list_count, deck_list_count, error_message, duration_seconds) VALUES (CURRENT_TIMESTAMP, :status, :cards, :sets, :identifiers, :card_types, :keywords, :sets_list, :deck_list, :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), 'sets_list': counts.get('set_list', 0), 'deck_list': counts.get('deck_list', 0), 'error': error, 'duration': duration, }) async def sync_mirrors(self): """Sync mirrored cards from mtg_cards to mtg_cards_mirror. This should be called after a successful refresh to ensure the mirror table is up-to-date with the latest card data. """ from sqlalchemy import insert, update logger.info("Syncing card mirrors...") async with mtg_async_session() as session: # Get all cards from mtg_cards cards_stmt = text(""" SELECT id, name, mana_cost, type_line, oracle_text, power, toughness, rarity, layout, artist, flavor_text, numbers, identifiers, images, image, card_parts, keywords, legalities, set_code, set_name FROM mtg_cards """) result = await session.execute(cards_stmt) cards = result.fetchall() synced_count = 0 for card in cards: # Upsert into mtg_cards_mirror await session.execute(text(""" INSERT INTO mtg_cards_mirror (source_id, name, mana_cost, type_line, oracle_text, power, toughness, rarity, layout, artist, flavor_text, numbers, identifiers, images, image, card_parts, keywords, legalities, set_code, set_name) VALUES (:source_id, :name, :mana_cost, :type_line, :oracle_text, :power, :toughness, :rarity, :layout, :artist, :flavor_text, :numbers, :identifiers, :images, :image, :card_parts, :keywords, :legalities, :set_code, :set_name) ON CONFLICT (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, image = EXCLUDED.image, card_parts = EXCLUDED.card_parts, keywords = EXCLUDED.keywords, legalities = EXCLUDED.legalities, set_code = EXCLUDED.set_code, set_name = EXCLUDED.set_name, synced_at = CURRENT_TIMESTAMP """), { 'source_id': card[0], 'name': card[1], 'mana_cost': card[2], 'type_line': card[3], 'oracle_text': card[4], 'power': card[5], 'toughness': card[6], 'rarity': card[7], 'layout': card[8], 'artist': card[9], 'flavor_text': card[10], 'numbers': card[11], 'identifiers': card[12], 'images': card[13], 'image': card[14], 'card_parts': card[15], 'keywords': card[16], 'legalities': card[17], 'set_code': card[18], 'set_name': card[19], }) synced_count += 1 await session.commit() logger.info(f"Synced {synced_count} card mirrors") return synced_count async def run_refresh(self) -> bool: """Run complete refresh cycle from local files.""" logger.info("Starting local file refresh") try: import time start_time = time.time() # Upsert data counts = await self.upsert_all() # Log refresh duration = int(time.time() - start_time) status = "SUCCESS" await self.log_refresh(status, counts, duration) logger.info(f"Refresh complete in {duration}s: {counts}") return True except Exception as e: logger.error(f"Refresh failed: {e}") await self.log_refresh("FAILED", {}, 0, str(e)) return False async def download_files(self, force: bool = False, max_retries: int = 3) -> dict: """ Download MTGJSON files from the API. Args: force: Force re-download even if files exist max_retries: Maximum number of retry attempts Returns: Dictionary with download results """ logger.info("Starting file download from MTGJSON API") # Check which files need downloading files_to_download = [] for filename in self._get_file_urls().keys(): file_path = self.data_dir / filename if force or not file_path.exists(): files_to_download.append(filename) else: logger.info(f"Skipping {filename} (already exists)") if not files_to_download: logger.info("All files already downloaded") return {"success": True, "downloaded": [], "skipped": list(self._get_file_urls().keys())} # Download files results = {} async with aiohttp.ClientSession() as session: for filename in files_to_download: success, message = await self._download_single_file( session, filename, max_retries ) results[filename] = {"success": success, "message": message} return results async def _download_single_file( self, session: aiohttp.ClientSession, filename: str, max_retries: int ) -> Tuple[bool, str]: """Download a single file with retry logic.""" url = self._get_file_urls()[filename]["url"] decompress = self._get_file_urls()[filename]["decompress"] for attempt in range(max_retries): try: logger.info(f"Downloading {filename} (attempt {attempt + 1}/{max_retries})") async with session.get(url, timeout=aiohttp.ClientTimeout(total=3600)) as response: if response.status != 200: error_msg = f"HTTP {response.status} for {filename}" logger.error(error_msg) if attempt < max_retries - 1: await asyncio.sleep(5) continue return False, error_msg # Write to temp file temp_path = self.data_dir / f".{filename}.tmp" with open(temp_path, 'wb') as f: async for chunk in response.content.iter_chunked(8192): f.write(chunk) # Move to final location final_path = self.data_dir / filename temp_path.rename(final_path) # Decompress if needed if decompress: await self._decompress_file(final_path) logger.info(f"Successfully downloaded {filename}") return True, "Download complete" except Exception as e: error_msg = f"Error downloading {filename}: {str(e)}" logger.error(error_msg) if attempt < max_retries - 1: await asyncio.sleep(5) continue return False, error_msg return False, f"Failed after {max_retries} attempts" def _get_file_urls(self) -> dict: """Get URL configuration for MTGJSON files.""" return { "AllPrintings.psql": { "url": "https://mtgjson.com/api/v5/AllPrintings.psql", "decompress": False, "description": "AllPrintings.psql" }, "AllIdentifiers.json": { "url": "https://mtgjson.com/api/v5/AllIdentifiers.json.gz", "decompress": True, "description": "AllIdentifiers.json" }, "Keywords.json": { "url": "https://mtgjson.com/api/v5/Keywords.json.gz", "decompress": True, "description": "Keywords.json" }, "CardTypes.json": { "url": "https://mtgjson.com/api/v5/CardTypes.json.gz", "decompress": True, "description": "CardTypes.json" }, "AllDeckFiles.zip": { "url": "https://mtgjson.com/api/v5/AllDeckFiles.zip", "decompress": False, "description": "AllDeckFiles.zip" }, } async def _decompress_file(self, filepath: Path) -> None: """Decompress a gzip file, detecting the target extension from URL config.""" logger.info(f"Decompressing {filepath.name}") # Determine target extension based on filename if filepath.name.endswith('.json.gz'): target_path = filepath.with_suffix('.json') elif filepath.name.endswith('.psql.gz'): target_path = filepath.with_suffix('.psql') else: # Default: remove .gz extension target_path = filepath.with_suffix('') with gzip.open(filepath, 'rb') as f_in: with open(target_path, 'wb') as f_out: f_out.write(f_in.read()) filepath.unlink() target_path.rename(filepath) logger.info(f"Decompressed {filepath.name} -> {target_path.name}") async def verify_files(self) -> Tuple[bool, List[str]]: """ Verify that all required files exist and are valid. Returns: Tuple of (all_valid, list_of_errors) """ errors = [] for filename in self._get_file_urls().keys(): file_path = self.data_dir / filename if not file_path.exists(): errors.append(f"Missing required file: {filename}") continue # Check file size (basic sanity check) size = file_path.stat().st_size if size == 0: errors.append(f"Empty file: {filename}") continue # Verify JSON files if filename.endswith('.json'): try: with open(file_path, 'r') as f: json.load(f) except Exception as e: errors.append(f"Invalid JSON in {filename}: {str(e)}") # Verify ZIP files if filename.endswith('.zip'): try: with zipfile.ZipFile(file_path, 'r') as zf: zf.testzip() except Exception as e: errors.append(f"Invalid ZIP file {filename}: {str(e)}") # Verify SQL files if filename.endswith('.psql'): with open(file_path, 'r') as f: content = f.read(1024) if not any(kw in content.upper() for kw in ['INSERT', 'CREATE', 'BEGIN']): errors.append(f"File {filename} doesn't appear to contain SQL") all_valid = len(errors) == 0 return all_valid, errors async def download_and_refresh(self, force: bool = False) -> dict: """ Download files and upsert data in one operation. Args: force: Force re-download Returns: Dictionary with download and upsert results """ import time start_time = time.time() result = { "download": None, "upsert": None, "success": False, "duration": None } try: # Download files logger.info("Starting download and refresh cycle") download_result = await self.download_files(force=force) result["download"] = download_result # Check if download was successful if not all(r["success"] for r in download_result.values()): failed_files = [f for f, r in download_result.items() if not r["success"]] error_msg = f"Download failed for: {', '.join(failed_files)}" logger.error(error_msg) await self.log_refresh("FAILED", {}, 0, error_msg) result["error"] = error_msg result["duration"] = int(time.time() - start_time) return result # Upsert data upsert_result = await self.upsert_all() result["upsert"] = upsert_result # Log refresh duration = int(time.time() - start_time) result["success"] = True result["duration"] = duration await self.log_refresh("SUCCESS", upsert_result, duration) logger.info(f"Download and refresh complete in {duration}s: {upsert_result}") return result except Exception as e: logger.error(f"Download and refresh failed: {e}") duration = int(time.time() - start_time) await self.log_refresh("FAILED", {}, duration, str(e)) result["error"] = str(e) result["duration"] = duration return result # Singleton instance _manager_instance = 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 if __name__ == "__main__": import argparse parser = argparse.ArgumentParser(description="MTGJSON Data Manager") parser.add_argument("--refresh", action="store_true", help="Run refresh cycle") parser.add_argument("--download", action="store_true", help="Download files") parser.add_argument("--force", action="store_true", help="Force re-download") parser.add_argument("--data-dir", type=str, default="/app/data", help="Data directory") args = parser.parse_args() async def main(): manager = MTGJSONManager(Path(args.data_dir)) if args.refresh: success = await manager.run_refresh() if success: print("Refresh completed successfully") else: print("Refresh failed") exit(1) elif args.download: result = await manager.download_files(force=args.force) print(f"Download result: {result}") else: print("Usage: python -m app.services.mtgjson_manager [--refresh | --download]") asyncio.run(main())