Files
image-db-sync-monitor/backend/core/scanner.py
2026-07-07 15:54:55 +05:30

207 lines
8.0 KiB
Python

import asyncio
import aiohttp
import aiomysql
import json
import os
from datetime import datetime, timezone
from urllib.parse import urljoin
DATA_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(__file__))), 'data')
CONFIG_PATH = os.path.join(DATA_DIR, 'config.json')
SCHEMA_RESULTS_PATH = os.path.join(DATA_DIR, 'schema_results.json')
IMAGE_RESULTS_PATH = os.path.join(DATA_DIR, 'image_results.json')
def load_config():
with open(CONFIG_PATH, 'r') as f:
return json.load(f)
def load_json(path):
if not os.path.exists(path):
return {}
try:
with open(path, 'r') as f:
return json.load(f)
except Exception:
return {}
def save_json(path, data):
with open(path, 'w') as f:
json.dump(data, f, indent=2)
SCAN_STATE = {
"schema_scanning": False,
"schema_current_db": None,
"image_scanning": False,
"image_current_db": None
}
def get_scan_status():
return SCAN_STATE
async def scan_schema_for_db(db_name: str, config: dict):
host = config['host']
port = config.get('port', 3306)
user = config['user']
password = config['password']
tables_data = []
try:
pool = await aiomysql.create_pool(
host=host, port=port, user=user, password=password, db=db_name, autocommit=True, connect_timeout=5
)
async with pool.acquire() as conn:
async with conn.cursor(aiomysql.DictCursor) as cur:
await cur.execute("SHOW TABLES")
tables = [list(row.values())[0] for row in await cur.fetchall()]
for table in tables:
try:
await cur.execute(f"SELECT COUNT(*) as cnt FROM `{table}`")
row_count = (await cur.fetchone())['cnt']
await cur.execute(f"SELECT COUNT(*) as cnt FROM information_schema.columns WHERE table_schema='{db_name}' AND table_name='{table}'")
col_count = (await cur.fetchone())['cnt']
tables_data.append({
"table_name": table,
"row_count": row_count,
"column_count": col_count,
"status": "Success"
})
except Exception as e:
tables_data.append({
"table_name": table,
"row_count": 0,
"column_count": 0,
"status": f"Error: {e}"
})
pool.close()
await pool.wait_closed()
return {"db_name": db_name, "status": "Success", "tables": tables_data}
except Exception as e:
return {"db_name": db_name, "status": f"Connection Failed: {e}", "tables": []}
async def check_url(session, url, semaphore, timeout_seconds):
async with semaphore:
try:
async with session.head(url, allow_redirects=True, timeout=timeout_seconds) as response:
if response.status == 405:
async with session.get(url, allow_redirects=True, timeout=timeout_seconds) as get_response:
return "Available" if get_response.status == 200 else "Not Found"
return "Available" if response.status == 200 else "Not Found"
except asyncio.TimeoutError:
return "Timeout"
except Exception:
return "Error"
async def scan_images_for_db_table(db_name: str, table: str, config: dict):
host = config['host']
port = config.get('port', 3306)
user = config['user']
password = config['password']
image_column = config.get('image_column', 'test_image_path')
base_url = config['base_url']
if not base_url.endswith('/'): base_url += '/'
concurrency_limit = config.get('concurrency_limit', 50)
timeout_seconds = config.get('request_timeout_seconds', 10)
stats = {"table": table, "total_rows_with_image": 0, "available": 0, "not_found": 0, "error": 0, "status": "Success"}
try:
pool = await aiomysql.create_pool(
host=host, port=port, user=user, password=password, db=db_name, autocommit=True, connect_timeout=5
)
async with pool.acquire() as conn:
async with conn.cursor(aiomysql.DictCursor) as cur:
try:
await cur.execute(f"SELECT `{image_column}` FROM `{table}` WHERE `{image_column}` IS NOT NULL AND `{image_column}` != ''")
rows = await cur.fetchall()
except Exception as e:
stats['status'] = f"Table/Query Error: {e}"
pool.close()
await pool.wait_closed()
return stats
pool.close()
await pool.wait_closed()
stats['total_rows_with_image'] = len(rows)
if len(rows) == 0: return stats
semaphore = asyncio.Semaphore(concurrency_limit)
connector = aiohttp.TCPConnector(limit=concurrency_limit)
async with aiohttp.ClientSession(connector=connector) as session:
tasks = []
for row in rows:
img_path = row[image_column]
if isinstance(img_path, str) and img_path.startswith('/'): img_path = img_path[1:]
url = urljoin(base_url, img_path)
tasks.append(asyncio.create_task(check_url(session, url, semaphore, timeout_seconds)))
for f in asyncio.as_completed(tasks):
res = await f
if res == "Available": stats['available'] += 1
elif res == "Not Found": stats['not_found'] += 1
else: stats['error'] += 1
return stats
except Exception as e:
stats['status'] = f"Connection/Execution Failed: {e}"
return stats
async def run_schema_scan():
if SCAN_STATE["schema_scanning"]: return
SCAN_STATE["schema_scanning"] = True
try:
config = load_config()
databases = config.get('databases', [])
results = load_json(SCHEMA_RESULTS_PATH)
for db_name in databases:
SCAN_STATE["schema_current_db"] = db_name
schema_data = await scan_schema_for_db(db_name, config)
results[db_name] = {
"schema_status": schema_data['status'],
"schema_tables": schema_data['tables'],
"last_updated": datetime.now(timezone.utc).isoformat()
}
save_json(SCHEMA_RESULTS_PATH, results) # Instantly available
finally:
SCAN_STATE["schema_scanning"] = False
SCAN_STATE["schema_current_db"] = None
async def run_image_scan():
if SCAN_STATE["image_scanning"]: return
SCAN_STATE["image_scanning"] = True
try:
config = load_config()
databases = config.get('databases', [])
image_tables = config.get('image_tables', [])
schema_results = load_json(SCHEMA_RESULTS_PATH)
image_results = load_json(IMAGE_RESULTS_PATH)
for db_name in databases:
SCAN_STATE["image_current_db"] = db_name
db_schema = schema_results.get(db_name, {})
if db_schema.get('schema_status') == "Success":
existing_tables = [t['table_name'] for t in db_schema.get('schema_tables', [])]
stats_list = []
for itable in image_tables:
if itable in existing_tables:
img_data = await scan_images_for_db_table(db_name, itable, config)
stats_list.append(img_data)
image_results[db_name] = {
"image_stats": stats_list,
"last_updated": datetime.now(timezone.utc).isoformat()
}
save_json(IMAGE_RESULTS_PATH, image_results) # Instantly available
finally:
SCAN_STATE["image_scanning"] = False
SCAN_STATE["image_current_db"] = None