diff --git a/api/main.py b/api/main.py index 9e558a3..348a569 100644 --- a/api/main.py +++ b/api/main.py @@ -1,69 +1,170 @@ """MedArchive API — эндпоинты ТЗ §4.5 плюс /stats для дашборда. -Здесь зафиксирован контракт API (пути, параметры, модели ответа). Реализация -запросов к БД — задача агента D; пока эндпоинты возвращают пустые заготовки, -чтобы приложение поднималось и фронт мог разрабатываться на сгенерированной -OpenAPI-схеме. +Читает базу, собранную конвейером (etl/pipeline.py). Для MVP это SQLite; путь +задаётся переменной окружения MEDARCHIVE_DB. Ядро WOW — /services/{id}/partners: +«кто оказывает услугу и по какой цене», отсортировано от выгодной. """ from __future__ import annotations +import json +import os +import sqlite3 import sys from pathlib import Path -from fastapi import FastAPI, Query +from fastapi import FastAPI, HTTPException, Query +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="0.1.0", + version="1.0.0", description="Поиск услуг и цен по архиву прайсов клиник-партнёров.", ) +app.add_middleware( + CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"] +) + + +def db() -> sqlite3.Connection: + conn = sqlite3.connect(DB_PATH) + conn.row_factory = sqlite3.Row + return conn + + +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 = Query(None, description="Фильтр по специальности")): - return [] # TODO(D): выборка из таблицы service +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: + sql += " AND name_ru LIKE ?" + args.append(f"%{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): - return [] # TODO(D): партнёры с ценами по услуге (ядро WOW-сравнения) + 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 + ORDER BY COALESCE(pi.price_resident, pi.price_nonresident)""" + with db() as conn: + return [_price_out(r) for r in conn.execute(sql, [service_id])] @app.get("/partners", response_model=list[PartnerOut], summary="Партнёры") -def list_partners(city: str | None = None, is_active: bool | None = None): - return [] # TODO(D) +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): - return [] # TODO(D) +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)): - return {"services": [], "partners": []} # TODO(D) +def search(q: str = Query(..., min_length=1), limit: int = 20): + with db() as conn: + services = [ + dict(r) + for r in conn.execute( + "SELECT service_id, specialty, name_ru FROM service WHERE name_ru LIKE ? LIMIT ?", + [f"%{q}%", limit], + ) + ] + partners = [ + dict(r) + for r in conn.execute( + "SELECT partner_id, name, city FROM partner WHERE name LIKE ? LIMIT ?", [f"%{q}%", limit] + ) + ] + return {"services": services, "partners": partners} @app.get("/unmatched", summary="Несопоставленные позиции (для операторов)") def unmatched(limit: int = 50, offset: int = 0): - return [] # TODO(D) + sql = """SELECT pi.item_id, pi.service_name_raw, pi.prices, p.name AS partner_name + FROM price_item pi JOIN partner p ON pi.partner_id = p.partner_id + WHERE pi.service_id IS NULL AND pi.is_active = 1 + ORDER BY 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"], + "prices": json.loads(r["prices"] or "{}"), + } + for r in conn.execute(sql, [limit, offset]) + ] @app.post("/match", summary="Ручное сопоставление позиции со справочником") -def manual_match(item_id: str, service_id: str): - return {"item_id": item_id, "service_id": service_id, "status": "todo"} # TODO(D) +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, "услуга справочника не найдена") + cur = conn.execute( + "UPDATE price_item SET service_id = ?, map_method = 'manual', map_confidence = 1.0 WHERE item_id = ?", + [service_id, item_id], + ) + conn.commit() + if cur.rowcount == 0: + raise HTTPException(404, "позиция не найдена") + return {"item_id": item_id, "service_id": service_id, "status": "matched"} @app.get("/stats", response_model=Stats, summary="Сводка качества обработки") def stats(): - return Stats() # TODO(D): реальные метрики из БД + 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"} + return {"status": "ok", "db": Path(DB_PATH).exists()} diff --git a/etl/dictionary.py b/etl/dictionary.py index 4046ca5..5933ff2 100644 --- a/etl/dictionary.py +++ b/etl/dictionary.py @@ -64,6 +64,7 @@ def load_dictionary(path: str | Path) -> list[Service]: raise ValueError(f"В справочнике нет колонки Name_ru. Заголовки: {header}") services: list[Service] = [] + seen_ids: dict[str, int] = {} for row in rows: raw_name = row[name_idx] if not raw_name or not str(raw_name).strip(): @@ -74,8 +75,15 @@ def load_dictionary(path: str | Path) -> list[Service]: if tar and not TARIFICATOR_RE.match(tar): tar = None # отбрасываем мусорные коды, оставляем только валидный формат + # «-» иногда повторяется в справочнике — гарантируем уникальность. + sid = f"{row[col['ID']]}-{row[col['Code']]}" + repeat = seen_ids.get(sid, 0) + seen_ids[sid] = repeat + 1 + if repeat: + sid = f"{sid}#{repeat}" + service = Service( - service_id=f"{row[col['ID']]}-{row[col['Code']]}", + service_id=sid, specialty=str(row[col["Специальность"]]).strip() if row[col["Специальность"]] else "", name_ru=str(raw_name).strip(), tarificator_code=tar, diff --git a/etl/extractors/__init__.py b/etl/extractors/__init__.py index 8a7696c..b0510f0 100644 --- a/etl/extractors/__init__.py +++ b/etl/extractors/__init__.py @@ -18,8 +18,12 @@ from etl.extractors import readers # noqa: E402 from etl.extractors.grid import extract_rows_from_grid # noqa: E402 -def extract(path: str) -> ParsedDocument: - """Разобрать один файл прайса в структурированный ParsedDocument.""" +def extract(path: str, use_vision: bool = True) -> ParsedDocument: + """Разобрать один файл прайса в структурированный ParsedDocument. + + use_vision=False отключает добор через Gemini Vision (быстрый прогон без + обращений к API), даже если ключ задан. + """ name = Path(path).name fmt = readers.detect_format(path) doc = ParsedDocument(file_name=name, file_format=fmt, partner_name=readers.partner_from_name(name)) @@ -55,7 +59,7 @@ def extract(path: str) -> ParsedDocument: doc.parse_log.append(f"ошибка разбора таблицы: {type(exc).__name__}: {exc}") # Трудные PDF (скан, битый слой, нулевой разбор) добираем через Gemini Vision. - needs_vision = fmt == "pdf" and (doc.file_format == "scan_pdf" or not doc.rows) + needs_vision = use_vision and fmt == "pdf" and (doc.file_format == "scan_pdf" or not doc.rows) if needs_vision and os.environ.get("GEMINI_API_KEY"): try: from etl.extractors.vision import extract_with_vision diff --git a/etl/pipeline.py b/etl/pipeline.py new file mode 100644 index 0000000..d026155 --- /dev/null +++ b/etl/pipeline.py @@ -0,0 +1,137 @@ +"""Конвейер обработки архива: извлечение → нормализация → валидация → SQLite. + +Для MVP данные складываются в SQLite (разворачивается без Docker), а боевая схема +PostgreSQL + pgvector лежит в db/migrations. Здесь же — дедупликация партнёров, +версионирование цен (последний прайс активен, старые архивируются) и сбор отчёта о +качестве для дашборда и сдачи. +""" +from __future__ import annotations + +import json +import sqlite3 +import sys +import uuid +from collections import Counter +from datetime import date +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parents[1])) + +from etl.dictionary import load_dictionary, normalize_name +from etl.extractors import extract +from etl.normalize import Matcher +from etl.validate import validate_item + +SCHEMA = """ +CREATE TABLE IF NOT EXISTS service ( + service_id TEXT PRIMARY KEY, specialty TEXT, name_ru TEXT, name_norm TEXT, tarificator_code TEXT +); +CREATE TABLE IF NOT EXISTS partner ( + partner_id TEXT PRIMARY KEY, name TEXT, name_norm TEXT UNIQUE, city TEXT +); +CREATE TABLE IF NOT EXISTS price_document ( + doc_id TEXT PRIMARY KEY, partner_id TEXT, file_name TEXT, file_format TEXT, + effective_date TEXT, parse_status TEXT, rows_count INTEGER +); +CREATE TABLE IF NOT EXISTS price_item ( + item_id INTEGER PRIMARY KEY AUTOINCREMENT, doc_id TEXT, partner_id TEXT, + service_name_raw TEXT, service_code_source TEXT, service_id TEXT, + prices TEXT, price_resident REAL, price_nonresident REAL, unit TEXT, + effective_date TEXT, map_method TEXT, map_confidence REAL, status TEXT, is_active INTEGER +); +CREATE INDEX IF NOT EXISTS price_item_service ON price_item(service_id); +CREATE INDEX IF NOT EXISTS price_item_partner ON price_item(partner_id); +""" + + +def _open(db_path: str) -> sqlite3.Connection: + conn = sqlite3.connect(db_path) + conn.executescript(SCHEMA) + for table in ("price_item", "price_document", "partner", "service"): + conn.execute(f"DELETE FROM {table}") # пересборка с нуля при каждом прогоне + return conn + + +def run(data_dir: str, db_path: str, dict_path: str, embedder=None, use_vision: bool = False, + today: date | None = None) -> dict: + """Прогнать весь архив в SQLite и вернуть отчёт о качестве.""" + today = today or date.today() + services = load_dictionary(dict_path) + matcher = Matcher(services, embedder=embedder) + + conn = _open(db_path) + conn.executemany( + "INSERT INTO service VALUES (?,?,?,?,?)", + [(s.service_id, s.specialty, s.name_ru, s.name_norm, s.tarificator_code) for s in services], + ) + + partners: dict[str, str] = {} + staged: list[tuple] = [] # (doc_id, partner_id, raw, code, prices, unit, eff_iso) + report: Counter = Counter() + + for path in sorted(Path(data_dir).glob("*")): + doc = extract(str(path), use_vision=use_vision) + partner_name = doc.partner_name or path.stem + partner_norm = normalize_name(partner_name) + partner_id = partners.get(partner_norm) + if partner_id is None: + partner_id = str(uuid.uuid4()) + partners[partner_norm] = partner_id + conn.execute("INSERT INTO partner VALUES (?,?,?,?)", (partner_id, partner_name, partner_norm, None)) + + doc_id = str(uuid.uuid4()) + eff_iso = doc.effective_date.isoformat() if doc.effective_date else None + conn.execute( + "INSERT INTO price_document VALUES (?,?,?,?,?,?,?)", + (doc_id, partner_id, doc.file_name, doc.file_format, eff_iso, "done", len(doc.rows)), + ) + report["documents"] += 1 + for row in doc.rows: + staged.append((doc_id, partner_id, row.service_name_raw, row.service_code_source, + row.prices, row.unit, eff_iso)) + + # Нормализация одной пачкой (эмбеддинги считаются разом — быстрее и дешевле). + matches = matcher.match_batch([s[2] for s in staged], [s[3] for s in staged]) + + for (doc_id, partner_id, raw, code, prices, unit, eff_iso), match in zip(staged, matches, strict=True): + status, _flags = validate_item(raw, prices, effective_date=None, today=today) + if match.service_id: + report["matched"] += 1 + report[f"method_{match.method}"] += 1 + conn.execute( + """INSERT INTO price_item + (doc_id, partner_id, service_name_raw, service_code_source, service_id, prices, + price_resident, price_nonresident, unit, effective_date, map_method, map_confidence, + status, is_active) + VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,1)""", + (doc_id, partner_id, raw, code, match.service_id, json.dumps(prices, ensure_ascii=False), + prices.get("resident"), prices.get("nonresident"), unit, eff_iso, match.method, + round(match.confidence, 3), status), + ) + report["items"] += 1 + + # Версионирование: если у партнёра по той же услуге есть более свежий прайс — старые архивируем. + conn.execute( + """UPDATE price_item SET is_active = 0 + WHERE effective_date IS NOT NULL AND EXISTS ( + SELECT 1 FROM price_item b + WHERE b.partner_id = price_item.partner_id + AND b.service_name_raw = price_item.service_name_raw + AND b.effective_date > price_item.effective_date)""" + ) + conn.commit() + + total = report["items"] + summary = { + "documents": report["documents"], + "partners": len(partners), + "items": total, + "auto_matched": report["matched"], + "auto_matched_pct": round(100 * report["matched"] / total, 1) if total else 0.0, + "unmatched": total - report["matched"], + "by_method": { + method: report[f"method_{method}"] for method in ("code", "exact", "embedding", "fuzzy", None) + }, + } + conn.close() + return summary diff --git a/scripts/run_pipeline.py b/scripts/run_pipeline.py new file mode 100644 index 0000000..b8168d6 --- /dev/null +++ b/scripts/run_pipeline.py @@ -0,0 +1,27 @@ +"""Запуск конвейера на архиве: data/raw → SQLite + отчёт о качестве. + +Флаг --vision включает добор трудных PDF через Gemini Vision (медленно). +""" +import json +import sys +from datetime import date +from pathlib import Path + +REPO = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(REPO)) + +from etl.normalize.embedding import GeminiEmbedder +from etl.pipeline import run + +summary = run( + data_dir=str(REPO / "data/raw"), + db_path=str(REPO / "data/medarchive.db"), + dict_path=str(REPO / "data/reference/dictionary.xlsx"), + embedder=GeminiEmbedder(), + use_vision="--vision" in sys.argv, + today=date(2026, 6, 26), +) +(REPO / "data/quality_report.json").write_text( + json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8" +) +print(json.dumps(summary, ensure_ascii=False, indent=2)) diff --git a/web/index.html b/web/index.html new file mode 100644 index 0000000..cb680c4 --- /dev/null +++ b/web/index.html @@ -0,0 +1,137 @@ + + + + + +MedArchive — сравнение цен на медуслуги + + + + + + + +
+ +
+
+ MedArchive + прайсы клиник-партнёров +
+ +
+ +
+ + + + + + + +
+ + + +