Files
medtech-hackathon/api/main.py
T
admins dbe36d3323 fix(normalize): confidence-aware guard переобобщения → честные 73,7%
Прежний guard снимал эмбеддинг-матч, если услуга собрала в клинике больше
8 разных названий. Это наказывало легитимные частотные услуги — одна «МРТ
головного мозга» реально встречается под десятками протоколов (с контрастом,
3Тл, предоперационная), и все они верно нормализуются к ней. Процент падал
до 69%.

Теперь снимаем только НЕуверенные матчи (<0,72) внутри концентрированных
групп: так отсекается сиблинг-захват (напр. «Магний Mg (моча)» подтягивал
Никель/Свинец/Ртуть в моче на пороге 0,70), а верные вариации остаются
сопоставленными. Признак ошибки — низкая уверенность, а не число названий.

Итог на тестовом архиве: 16 140 позиций извлечено, 73,7% нормализовано
автоматически (цель ТЗ ≥70%), 273 неуверенных матча — в очередь верификации.
Цифры выровнены в README, ARCHITECTURE, питч-деках, PPTX, quality_report.

Плюс полировка: извлечение города клиники, регистронезависимый поиск.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-27 19:42:21 +05:00

587 lines
27 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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, Response, 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)")
# Настройки админки (ключ-значение): сейчас — путь к папке авто-загрузки.
conn.execute("CREATE TABLE IF NOT EXISTS settings (key TEXT PRIMARY KEY, value TEXT)")
# Журнал уже обработанных файлов авто-загрузки (чтобы не обрабатывать повторно).
conn.execute(
"CREATE TABLE IF NOT EXISTS ingested_file ("
"name TEXT PRIMARY KEY, source TEXT, ingested_at TEXT, items INTEGER)"
)
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("\\", "/"))
def _ingest_dir(data_dir: str) -> tuple[dict, bool]:
"""Догрузить все прайсы из каталога в базу. Возвращает (отчёт, использованы_ли_эмбеддинги).
Тяжёлые зависимости (конвейер, Gemini) импортируются лениво — read-only API живёт и без
них. Эмбеддинги — по возможности: при сбое Gemini (квота/сеть) обрабатываем без них (fuzzy),
загрузка всё равно проходит.
"""
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=data_dir, db_path=DB_PATH, dict_path=dict_path)
try:
return ingest(embedder=embedder, **kwargs), embedder is not None
except Exception:
if embedder is None:
raise
return ingest(embedder=None, **kwargs), False
@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)"
)
summary, embeddings = _ingest_dir(str(extract_dir))
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": embeddings,
}
finally:
import shutil as _sh
_sh.rmtree(workdir, ignore_errors=True)
@app.get("/settings", summary="Настройки админки (папка авто-загрузки)")
def get_settings():
with db() as conn:
rows = {r["key"]: r["value"] for r in conn.execute("SELECT key, value FROM settings")}
return {"watch_folder": rows.get("watch_folder", "")}
@app.post("/settings", summary="Сохранить папку авто-загрузки")
def set_settings(watch_folder: str = ""):
with db() as conn:
conn.execute(
"INSERT INTO settings (key, value) VALUES ('watch_folder', ?) "
"ON CONFLICT(key) DO UPDATE SET value = excluded.value",
[watch_folder.strip()],
)
conn.commit()
return {"watch_folder": watch_folder.strip()}
@app.post("/scan", summary="Проверить папку авто-загрузки и обработать новые прайсы")
def scan_folder():
"""Сканирует папку из настроек, находит новые файлы и догружает их в базу.
Новизна определяется по имени файла (журнал `ingested_file`) — повторно не обрабатываем.
Эту же ручку дёргает крон: система забирает свежие прайсы из папки без участия человека.
"""
import shutil
import tempfile
from datetime import UTC, datetime
with db() as conn:
row = conn.execute("SELECT value FROM settings WHERE key = 'watch_folder'").fetchone()
folder = (row["value"] if row else "") or ""
done = {r["name"] for r in conn.execute("SELECT name FROM ingested_file")}
if not folder:
raise HTTPException(400, "папка авто-загрузки не задана — укажите её в настройках")
src = Path(folder)
if not src.is_dir():
raise HTTPException(404, f"папка не найдена: {folder}")
new_files = [
p
for p in sorted(src.iterdir())
if p.is_file() and p.suffix.lower() in _SUPPORTED_EXT and p.name not in done
]
if not new_files:
return {"status": "ok", "new_files": 0, "message": "новых прайсов нет"}
workdir = Path(tempfile.mkdtemp(prefix="medarchive_scan_"))
try:
for p in new_files:
shutil.copy(p, workdir / p.name)
summary, embeddings = _ingest_dir(str(workdir))
now = datetime.now(UTC).isoformat(timespec="seconds")
with db() as conn:
conn.executemany(
"INSERT OR REPLACE INTO ingested_file (name, source, ingested_at, items) "
"VALUES (?, 'folder', ?, NULL)",
[(p.name, now) for p in new_files],
)
conn.commit()
return {
"status": "ok",
"new_files": len(new_files),
"files": [p.name for p in new_files],
"items": summary["items"],
"auto_matched_pct": summary["auto_matched_pct"],
"embeddings": embeddings,
}
finally:
shutil.rmtree(workdir, ignore_errors=True)
@app.get("/price-changes", summary="Услуги, у которых цена изменилась между прайсами")
def price_changes(threshold: float = 10.0, limit: int = 100):
"""Сравнивает активную цену с предыдущей (архивной) у того же партнёра и услуги.
Возвращает позиции, где модуль изменения ≥ threshold% — детектор аномалий цен (ТЗ §4.4),
вынесенный в интерфейс. Предыдущая цена — самый свежий архивный прайс той же позиции.
"""
sql = """
SELECT p.name AS partner_name, a.service_name_raw, s.name_ru AS service_ru,
a.price_resident AS new_price, a.effective_date AS new_date,
(SELECT b.price_resident FROM price_item b
WHERE b.partner_id = a.partner_id
AND ((a.service_id IS NOT NULL AND b.service_id = a.service_id)
OR (a.service_id IS NULL AND b.service_name_raw = a.service_name_raw))
AND b.is_active = 0 AND b.price_resident IS NOT NULL
ORDER BY b.effective_date DESC LIMIT 1) AS old_price,
(SELECT b.effective_date FROM price_item b
WHERE b.partner_id = a.partner_id
AND ((a.service_id IS NOT NULL AND b.service_id = a.service_id)
OR (a.service_id IS NULL AND b.service_name_raw = a.service_name_raw))
AND b.is_active = 0 AND b.price_resident IS NOT NULL
ORDER BY b.effective_date DESC LIMIT 1) AS old_date
FROM price_item a
JOIN partner p ON a.partner_id = p.partner_id
LEFT JOIN service s ON a.service_id = s.service_id
WHERE a.is_active = 1 AND a.price_resident IS NOT NULL
"""
changes = []
with db() as conn:
for r in conn.execute(sql):
old = r["old_price"]
new = r["new_price"]
if not old:
continue
delta_pct = round((new - old) / old * 100, 1)
if abs(delta_pct) >= threshold:
changes.append(
{
"partner_name": r["partner_name"],
"service": r["service_ru"] or r["service_name_raw"],
"old_price": old,
"new_price": new,
"delta_pct": delta_pct,
"old_date": r["old_date"],
"new_date": r["new_date"],
}
)
changes.sort(key=lambda c: abs(c["delta_pct"]), reverse=True)
return {"threshold": threshold, "count": len(changes), "changes": changes[:limit]}
@app.get("/export.xlsx", summary="Выгрузить нормализованный сравнительный прайс в Excel")
def export_xlsx():
"""Единый каталог: услуга справочника → клиники и цены. Готовый файл для биллинга/закупок."""
from io import BytesIO
from openpyxl import Workbook
from openpyxl.styles import Font, PatternFill
# Один ряд на «услуга × клиника» — чистый сравнительный каталог без дублей сырых позиций.
sql = """SELECT s.name_ru AS service, s.specialty, p.name AS partner,
MIN(pi.price_resident) AS price_resident,
MIN(pi.price_nonresident) AS price_nonresident,
MAX(pi.effective_date) AS effective_date
FROM price_item pi
JOIN service s ON pi.service_id = s.service_id
JOIN partner p ON pi.partner_id = p.partner_id
WHERE pi.is_active = 1
GROUP BY s.service_id, p.partner_id
ORDER BY s.name_ru, price_resident"""
wb = Workbook()
ws = wb.active
ws.title = "Каталог услуг и цен"
headers = ["Услуга", "Специальность", "Клиника", "Цена резидент, ₸",
"Цена нерезидент, ₸", "Прайс от"]
ws.append(headers)
head_font = Font(bold=True, color="FFFFFF")
head_fill = PatternFill("solid", fgColor="FF4713")
for cell in ws[1]:
cell.font = head_font
cell.fill = head_fill
with db() as conn:
for r in conn.execute(sql):
ws.append([r["service"], r["specialty"], r["partner"], r["price_resident"],
r["price_nonresident"], r["effective_date"]])
widths = [46, 22, 26, 16, 18, 12]
for i, w in enumerate(widths, 1):
ws.column_dimensions[ws.cell(row=1, column=i).column_letter].width = w
ws.freeze_panes = "A2"
# Второй лист — ВСЕ активные позиции как есть (включая очередь верификации),
# чтобы выгрузка была полной: первый лист — чистый каталог, второй — всё сырьё.
ws2 = wb.create_sheet("Все позиции")
ws2.append(["Услуга (как в прайсе)", "Сопоставлено со справочником", "Специальность",
"Клиника", "Цена резидент, ₸", "Цена нерезидент, ₸", "Прайс от", "Статус"])
for cell in ws2[1]:
cell.font = head_font
cell.fill = head_fill
sql_all = """SELECT pi.service_name_raw, s.name_ru AS matched, s.specialty, p.name AS partner,
pi.price_resident, pi.price_nonresident, pi.effective_date,
CASE WHEN pi.service_id IS NULL THEN 'в очереди' ELSE 'сопоставлено' END st
FROM price_item pi
JOIN partner p ON pi.partner_id = p.partner_id
LEFT JOIN service s ON pi.service_id = s.service_id
WHERE pi.is_active = 1
ORDER BY pi.service_name_raw"""
with db() as conn:
for r in conn.execute(sql_all):
ws2.append([r["service_name_raw"], r["matched"], r["specialty"], r["partner"],
r["price_resident"], r["price_nonresident"], r["effective_date"], r["st"]])
for i, w in enumerate([40, 40, 20, 24, 16, 18, 12, 14], 1):
ws2.column_dimensions[ws2.cell(row=1, column=i).column_letter].width = w
ws2.freeze_panes = "A2"
buffer = BytesIO()
wb.save(buffer)
return Response(
content=buffer.getvalue(),
media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
headers={"Content-Disposition": 'attachment; filename="medarchive_prices.xlsx"'},
)
@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").fetchone()[0]
matched = conn.execute(
"SELECT COUNT(*) FROM price_item WHERE 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()}