"""MedArchive API — эндпоинты ТЗ §4.5 плюс /stats для дашборда. Читает базу, собранную конвейером (etl/pipeline.py). Для MVP это SQLite; путь задаётся переменной окружения MEDARCHIVE_DB. Ядро WOW — /services/{id}/partners: «кто оказывает услугу и по какой цене», отсортировано от выгодной. """ from __future__ import annotations import json import os import re import sqlite3 import sys from pathlib import Path from fastapi import FastAPI, File, HTTPException, Query, UploadFile from fastapi.middleware.cors import CORSMiddleware sys.path.insert(0, str(Path(__file__).resolve().parents[1])) from contracts.models import PartnerOut, PriceOut, ServiceOut, Stats # noqa: E402 DB_PATH = os.environ.get( "MEDARCHIVE_DB", str(Path(__file__).resolve().parents[1] / "data/medarchive.db") ) app = FastAPI( title="MedArchive API", version="1.0.0", description="Поиск услуг и цен по архиву прайсов клиник-партнёров.", # За Caddy сервис живёт на /api — без этого Swagger ищет спеку в корне и падает. root_path="/api", ) app.add_middleware(CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"]) def db() -> sqlite3.Connection: conn = sqlite3.connect(DB_PATH) conn.row_factory = sqlite3.Row # Отложенные оператором позиции (чтобы пропущенное не возвращалось в начало очереди). conn.execute("CREATE TABLE IF NOT EXISTS skipped_item (item_id INTEGER PRIMARY KEY)") return conn _WS = re.compile(r"\s+") def _norm(text: str | None) -> str: """Нормализованная форма названия (для запоминания синонима при верификации).""" return _WS.sub(" ", (text or "").strip().lower().replace("ё", "е").replace("-", " ")).strip() def _price_out(row: sqlite3.Row) -> PriceOut: return PriceOut( partner_id=row["partner_id"], partner_name=row["partner_name"], price_resident_kzt=row["price_resident"], price_nonresident_kzt=row["price_nonresident"], prices=json.loads(row["prices"] or "{}"), effective_date=row["effective_date"], ) @app.get("/services", response_model=list[ServiceOut], summary="Услуги справочника") def list_services(specialty: str | None = None, q: str | None = None, limit: int = 100): sql = "SELECT service_id, specialty, name_ru, tarificator_code FROM service WHERE 1=1" args: list = [] if specialty: sql += " AND specialty = ?" args.append(specialty) if q: # Поиск по нормализованной колонке: SQLite LIKE сворачивает регистр только для # латиницы, поэтому «УЗИ» и «узи» совпадут лишь при сравнении lowercase-форм. sql += " AND name_norm LIKE ?" args.append(f"%{_norm(q)}%") sql += " ORDER BY name_ru LIMIT ?" args.append(limit) with db() as conn: return [ServiceOut(**dict(r)) for r in conn.execute(sql, args)] @app.get( "/services/{service_id}/partners", response_model=list[PriceOut], summary="Кто оказывает услугу и по какой цене", ) def service_partners(service_id: str): """Одна строка на клинику с показательной ценой (резидент или минимальный тариф).""" sql = """SELECT pi.partner_id, p.name AS partner_name, pi.price_resident, pi.price_nonresident, pi.prices, pi.effective_date FROM price_item pi JOIN partner p ON pi.partner_id = p.partner_id WHERE pi.service_id = ? AND pi.is_active = 1""" best: dict[str, PriceOut] = {} with db() as conn: for r in conn.execute(sql, [service_id]): prices = json.loads(r["prices"] or "{}") display = r["price_resident"] or (min(prices.values()) if prices else None) current = best.get(r["partner_id"]) if current is None or ( display is not None and (current.price_resident_kzt is None or display < current.price_resident_kzt) ): best[r["partner_id"]] = PriceOut( partner_id=r["partner_id"], partner_name=r["partner_name"], price_resident_kzt=display, price_nonresident_kzt=r["price_nonresident"], prices=prices, effective_date=r["effective_date"], ) return sorted(best.values(), key=lambda x: x.price_resident_kzt or float("inf")) @app.get("/partners", response_model=list[PartnerOut], summary="Партнёры") def list_partners(city: str | None = None): sql = "SELECT partner_id, name, city FROM partner WHERE 1=1" args: list = [] if city: sql += " AND city = ?" args.append(city) sql += " ORDER BY name" with db() as conn: return [ PartnerOut(partner_id=r["partner_id"], name=r["name"], city=r["city"]) for r in conn.execute(sql, args) ] @app.get( "/partners/{partner_id}/services", response_model=list[PriceOut], summary="Все услуги партнёра с ценами", ) def partner_services(partner_id: str, limit: int = 500): sql = """SELECT pi.partner_id, p.name AS partner_name, pi.price_resident, pi.price_nonresident, pi.prices, pi.effective_date FROM price_item pi JOIN partner p ON pi.partner_id = p.partner_id WHERE pi.partner_id = ? AND pi.is_active = 1 LIMIT ?""" with db() as conn: return [_price_out(r) for r in conn.execute(sql, [partner_id, limit])] @app.get("/search", summary="Полнотекстовый поиск по услугам и партнёрам") def search(q: str = Query(..., min_length=1), limit: int = 20): # Сравниваем по name_norm (lowercase): иначе кириллический регистр («УЗИ» vs «узи») # не сворачивается и часть запросов даёт «ничего не найдено». like = f"%{_norm(q)}%" with db() as conn: services = [ dict(r) for r in conn.execute( "SELECT service_id, specialty, name_ru FROM service WHERE name_norm LIKE ? LIMIT ?", [like, limit], ) ] partners = [ dict(r) for r in conn.execute( "SELECT partner_id, name, city FROM partner WHERE name_norm LIKE ? LIMIT ?", [like, limit], ) ] return {"services": services, "partners": partners} @app.get("/unmatched", summary="Очередь верификации: позиции с предложенным кандидатом") def unmatched(limit: int = 50, offset: int = 0): """Несопоставленные позиции с лучшим кандидатом из справочника (сортировка — по уверенности).""" sql = """SELECT pi.item_id, pi.service_name_raw, pi.prices, pi.unit, p.name AS partner_name, pi.suggested_service_id, pi.suggested_score, s.name_ru AS suggested_name, s.specialty AS suggested_specialty FROM price_item pi JOIN partner p ON pi.partner_id = p.partner_id LEFT JOIN service s ON pi.suggested_service_id = s.service_id WHERE pi.service_id IS NULL AND pi.is_active = 1 AND pi.item_id NOT IN (SELECT item_id FROM skipped_item) ORDER BY pi.suggested_score DESC, pi.item_id LIMIT ? OFFSET ?""" with db() as conn: return [ { "item_id": r["item_id"], "service_name_raw": r["service_name_raw"], "partner_name": r["partner_name"], "unit": r["unit"], "prices": json.loads(r["prices"] or "{}"), "suggested": ( { "service_id": r["suggested_service_id"], "name_ru": r["suggested_name"], "specialty": r["suggested_specialty"], "score": r["suggested_score"], } if r["suggested_service_id"] else None ), } for r in conn.execute(sql, [limit, offset]) ] @app.get("/unmatched/count", summary="Сколько позиций ещё в очереди верификации") def unmatched_count(): with db() as conn: n = conn.execute( """SELECT COUNT(*) FROM price_item WHERE service_id IS NULL AND is_active = 1 AND item_id NOT IN (SELECT item_id FROM skipped_item)""" ).fetchone()[0] learned = conn.execute("SELECT COUNT(*) FROM learned_synonym").fetchone()[0] return {"remaining": n, "learned_synonyms": learned} @app.post("/skip", summary="Отложить позицию (пропустить в очереди верификации)") def skip_item(item_id: int): with db() as conn: conn.execute("INSERT OR IGNORE INTO skipped_item (item_id) VALUES (?)", [item_id]) conn.commit() return {"item_id": item_id, "status": "skipped"} _SUPPORTED_EXT = (".pdf", ".xlsx", ".xls", ".docx") def _safe_member_name(member) -> str: """Имя файла из ZIP: чинит кириллицу (cp866) и срезает путь (защита от traversal).""" name = member.filename if not (member.flag_bits & 0x800): # нет UTF-8-флага — вероятно cp437/cp866 (рус. Windows) try: name = name.encode("cp437").decode("cp866") except (UnicodeEncodeError, UnicodeDecodeError): pass return os.path.basename(name.replace("\\", "/")) @app.post("/ingest", summary="Загрузить ZIP-архив прайсов и обработать (догрузка в базу)") def ingest_archive(file: UploadFile = File(...)): """Приём архива через интерфейс (ТЗ §4.1): распаковать ZIP и догрузить прайсы в базу. Тяжёлые зависимости (конвейер, Gemini) импортируются лениво — read-only API живёт и без них. Эмбеддинги — по возможности: при сбое Gemini обрабатываем без них (fuzzy), архив всё равно принимается. """ import shutil import tempfile import zipfile name = (file.filename or "").lower() if not name.endswith(".zip"): raise HTTPException(400, "ожидается ZIP-архив (.zip)") workdir = Path(tempfile.mkdtemp(prefix="medarchive_ingest_")) try: zip_path = workdir / "upload.zip" with zip_path.open("wb") as out: shutil.copyfileobj(file.file, out) extract_dir = workdir / "files" extract_dir.mkdir() try: with zipfile.ZipFile(zip_path) as zf: members = [m for m in zf.infolist() if not m.is_dir()] for member in members: fname = _safe_member_name(member) if not fname or not fname.lower().endswith(_SUPPORTED_EXT): continue with zf.open(member) as src, (extract_dir / fname).open("wb") as dst: shutil.copyfileobj(src, dst) except zipfile.BadZipFile as exc: raise HTTPException(400, "файл не является корректным ZIP-архивом") from exc prices = [p for p in extract_dir.iterdir() if p.is_file()] if not prices: raise HTTPException( 422, "в архиве не найдено прайсов поддерживаемых форматов (pdf, xlsx, xls, docx)" ) from etl.pipeline import ingest # noqa: PLC0415 — ленивый импорт тяжёлых зависимостей embedder = None if os.environ.get("GEMINI_API_KEY"): try: from etl.normalize.embedding import GeminiEmbedder # noqa: PLC0415 embedder = GeminiEmbedder() except Exception: # ключа/SDK нет — нормализуем без эмбеддингов embedder = None dict_path = str(Path(DB_PATH).parent / "reference/dictionary.xlsx") kwargs = dict(data_dir=str(extract_dir), db_path=DB_PATH, dict_path=dict_path) try: summary = ingest(embedder=embedder, **kwargs) degraded = False except Exception: # Сбой эмбеддингов (квота/сеть) не должен ронять приём архива — обрабатываем без них. if embedder is None: raise summary = ingest(embedder=None, **kwargs) degraded = True return { "status": "ok", "files_accepted": len(prices), "documents": summary["documents"], "new_partners": summary["partners"], "items": summary["items"], "auto_matched": summary["auto_matched"], "auto_matched_pct": summary["auto_matched_pct"], "by_method": summary["by_method"], "embeddings": not degraded and embedder is not None, } finally: import shutil as _sh _sh.rmtree(workdir, ignore_errors=True) @app.post("/match", summary="Ручное сопоставление позиции (с дообучением)") def manual_match(item_id: int, service_id: str): with db() as conn: if not conn.execute("SELECT 1 FROM service WHERE service_id = ?", [service_id]).fetchone(): raise HTTPException(404, "услуга справочника не найдена") row = conn.execute( "SELECT service_name_raw FROM price_item WHERE item_id = ?", [item_id] ).fetchone() if not row: raise HTTPException(404, "позиция не найдена") conn.execute( "UPDATE price_item SET service_id = ?, map_method = 'manual', map_confidence = 1.0 " "WHERE item_id = ?", [service_id, item_id], ) # Дообучение: запоминаем синоним — следующие прогоны сопоставят его автоматически. conn.execute( "INSERT OR REPLACE INTO learned_synonym (name_norm, service_id) VALUES (?, ?)", [_norm(row["service_name_raw"]), service_id], ) conn.commit() return {"item_id": item_id, "service_id": service_id, "status": "matched", "learned": True} @app.get("/stats", response_model=Stats, summary="Сводка качества обработки") def stats(): """Считается из базы — так после загрузки нового архива цифры обновляются сразу.""" with db() as conn: docs = conn.execute("SELECT COUNT(*) FROM price_document").fetchone()[0] items = conn.execute("SELECT COUNT(*) FROM price_item WHERE is_active = 1").fetchone()[0] matched = conn.execute( "SELECT COUNT(*) FROM price_item WHERE is_active = 1 AND service_id IS NOT NULL" ).fetchone()[0] if items: return Stats( documents_total=docs, documents_done=docs, items_total=items, auto_matched_pct=round(100 * matched / items, 1), unmatched_total=items - matched, ) # Фолбэк на отчёт сборки, если база ещё пуста. report = Path(DB_PATH).parent / "quality_report.json" if report.exists(): data = json.loads(report.read_text(encoding="utf-8")) return Stats( documents_total=data.get("documents", 0), documents_done=data.get("documents", 0), items_total=data.get("items", 0), auto_matched_pct=data.get("auto_matched_pct", 0.0), unmatched_total=data.get("unmatched", 0), ) return Stats() @app.get("/health", summary="Проверка живости") def health(): return {"status": "ok", "db": Path(DB_PATH).exists()}