feat: авто-загрузка из папки, алерты об изменении цен, экспорт каталога в Excel

- POST /scan: сервис забирает новые прайсы из папки (настройка watch_folder),
  журнал ingested_file исключает повторную обработку; вызывается кроном.
- GET /price-changes: услуги, где цена изменилась не меньше порога между прайсами
  (детектор аномалий из ТЗ §4.4, вынесенный в интерфейс).
- GET /export.xlsx: единый сравнительный каталог услуг и цен (openpyxl).
- GET/POST /settings: папка авто-загрузки.
- админка: карточки авто-загрузки и изменений цен + кнопка экспорта.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-06-27 13:13:38 +05:00
parent 195b8fbc5e
commit 51df0aa490
2 changed files with 278 additions and 27 deletions
+206 -25
View File
@@ -14,7 +14,7 @@ import sqlite3
import sys
from pathlib import Path
from fastapi import FastAPI, File, HTTPException, Query, UploadFile
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]))
@@ -39,6 +39,13 @@ def db() -> sqlite3.Connection:
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
@@ -233,6 +240,34 @@ def _safe_member_name(member) -> str:
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 и догрузить прайсы в базу.
@@ -275,29 +310,7 @@ def ingest_archive(file: UploadFile = File(...)):
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
summary, embeddings = _ingest_dir(str(extract_dir))
return {
"status": "ok",
"files_accepted": len(prices),
@@ -307,7 +320,7 @@ def ingest_archive(file: UploadFile = File(...)):
"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,
"embeddings": embeddings,
}
finally:
import shutil as _sh
@@ -315,6 +328,174 @@ def ingest_archive(file: UploadFile = File(...)):
_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 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 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"
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: