"""Конвейер обработки архива: извлечение → нормализация → валидация → SQLite. Для MVP данные складываются в SQLite (разворачивается без Docker), а боевая схема PostgreSQL + pgvector лежит в db/migrations. Здесь же — дедупликация партнёров, версионирование цен (последний прайс активен, старые архивируются) и сбор отчёта о качестве для дашборда и сдачи. Два входа: `run()` собирает базу с нуля (полная пересборка), `ingest()` догружает новые прайсы в уже собранную базу (append) — на нём держится приём ZIP через интерфейс. """ from __future__ import annotations import json import sqlite3 import sys import uuid from collections import Counter from datetime import date from pathlib import Path import numpy as np 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, suggested_service_id TEXT, suggested_score REAL ); -- Синонимы, выученные при ручной верификации (для дообучения нормализации). CREATE TABLE IF NOT EXISTS learned_synonym (name_norm TEXT PRIMARY KEY, service_id TEXT); 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); """ _INSERT_ITEM = """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, suggested_service_id, suggested_score) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,1,?,?)""" def _emb_path(db_path: str) -> Path: """Путь к кэшу эмбеддингов справочника рядом с базой.""" return Path(db_path).parent / "dict_emb.npy" def _open(db_path: str) -> sqlite3.Connection: # Полная пересборка: удаляем файл, чтобы схема всегда создавалась свежей. Path(db_path).unlink(missing_ok=True) conn = sqlite3.connect(db_path) conn.row_factory = sqlite3.Row conn.executescript(SCHEMA) return conn def _load_services(conn, services) -> None: """Залить справочник услуг, если таблица ещё пуста.""" if conn.execute("SELECT 1 FROM service LIMIT 1").fetchone(): return conn.executemany( "INSERT INTO service VALUES (?,?,?,?,?)", [(s.service_id, s.specialty, s.name_ru, s.name_norm, s.tarificator_code) for s in services], ) def _process_documents(conn, data_dir, matcher, partners, today, use_vision) -> Counter: """Извлечь файлы из каталога, сопоставить со справочником и записать позиции. Мутирует `conn` и словарь `partners` (норм. имя → partner_id, чтобы один и тот же партнёр из разных файлов не задвоился). Возвращает счётчик для отчёта о качестве. """ report: Counter = Counter() staged: list[tuple] = [] # (doc_id, partner_id, raw, code, prices, unit, eff_iso) for path in sorted(Path(data_dir).glob("*")): if path.is_dir(): continue 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), ) report["new_partners"] += 1 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_ITEM, ( 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, match.suggested_service_id, round(match.suggested_score, 3) if match.suggested_score else None, ), ) report["items"] += 1 return report def _apply_versioning(conn) -> None: """У партнёра по той же услуге более свежий прайс вытесняет старые (is_active = 0).""" 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)""" ) def _summary(report: Counter) -> dict: total = report["items"] return { "documents": report["documents"], "partners": report["new_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) }, } def run( data_dir: str, db_path: str, dict_path: str, embedder=None, use_vision: bool = False, today: date | None = None, ) -> dict: """Собрать базу с нуля по всему архиву и вернуть отчёт о качестве.""" today = today or date.today() services = load_dictionary(dict_path) matcher = Matcher(services, embedder=embedder) # Кэшируем эмбеддинги справочника: при последующих догрузках их не пересчитываем. if matcher.dict_emb is not None: np.save(_emb_path(db_path), matcher.dict_emb) conn = _open(db_path) _load_services(conn, services) partners: dict[str, str] = {} report = _process_documents(conn, data_dir, matcher, partners, today, use_vision) _apply_versioning(conn) conn.commit() summary = _summary(report) conn.close() return summary def ingest( data_dir: str, db_path: str, dict_path: str, embedder=None, use_vision: bool | None = None, today: date | None = None, ) -> dict: """Догрузить новые прайсы из каталога в существующую базу (append, без пересборки). На этом держится приём ZIP через интерфейс: справочник и его эмбеддинги уже готовы, поэтому считается только новый файл. Партнёры дедуплицируются с уже загруженными, версионирование вытесняет устаревшие прайсы той же клиники. """ today = today or date.today() services = load_dictionary(dict_path) emb_path = _emb_path(db_path) dict_emb = np.load(emb_path) if (embedder is not None and emb_path.exists()) else None matcher = Matcher(services, embedder=embedder, dict_emb=dict_emb) if use_vision is None: use_vision = embedder is not None # есть ключ Gemini — разрешаем распознавание сканов conn = sqlite3.connect(db_path) conn.row_factory = sqlite3.Row conn.executescript(SCHEMA) _load_services(conn, services) # Уже загруженные партнёры — чтобы повторная выгрузка той же клиники не задвоила её. partners = { row["name_norm"]: row["partner_id"] for row in conn.execute("SELECT partner_id, name_norm FROM partner") } report = _process_documents(conn, data_dir, matcher, partners, today, use_vision) _apply_versioning(conn) conn.commit() summary = _summary(report) conn.close() return summary