Files
site/WESP_REL/app/lab/etl/reference_db_import.py
2026-07-17 12:57:18 +03:00

266 lines
8.4 KiB
Python

"""Import tab reference PostgreSQL → WESP SQLite (offline admin ETL)."""
from __future__ import annotations
import json
import os
from dataclasses import dataclass, field
from datetime import datetime
from typing import Any
import psycopg2
from psycopg2.extras import RealDictCursor
from app import db
from app.lab.models import LabAnimalProfile
from app.models import Component, Recipe, WESP_SUPPRESS_SYNC_ENQUEUE
from app.lab.services.profile_norms import import_norms_from_legacy_text, save_norms_from_payload
DEFAULT_PG_URL = "postgresql://neoton:neoton_secret@localhost:5432/neoton"
@dataclass
class ImportStats:
components_enriched: int = 0
components_unmatched: int = 0
tab_components_purged: int = 0
animal_profiles: int = 0
recipe_links: int = 0
errors: list[str] = field(default_factory=list)
def _pg_url() -> str:
return os.environ.get("TAB_REFERENCE_DATABASE_URL", DEFAULT_PG_URL).strip()
def _json_text(value: Any) -> str:
if value is None:
return "{}"
if isinstance(value, str):
return value if value else "{}"
return json.dumps(value, ensure_ascii=False)
def _dt(value: Any) -> datetime | None:
if value is None:
return None
if isinstance(value, datetime):
return value
return None
def _fetch_all(conn, sql: str) -> list[dict[str, Any]]:
with conn.cursor(cursor_factory=RealDictCursor) as cur:
cur.execute(sql)
return list(cur.fetchall())
def _normalize_name(name: str) -> str:
return " ".join((name or "").lower().split())
def _find_wesp_component(row: dict[str, Any]) -> Component | None:
"""Match tab feed_ingredient → existing WESP component (never create)."""
external_no = row.get("external_no")
name = (row.get("name") or "").strip()
norm = _normalize_name(name)
if external_no is not None:
hit = Component.query.filter_by(external_no=external_no, is_deleted=False).first()
if hit is not None:
return hit
if name:
hit = Component.query.filter(
Component.name == name, Component.is_deleted.is_(False)
).first()
if hit is not None:
return hit
if norm:
for comp in Component.query.filter(Component.is_deleted.is_(False)).all():
if _normalize_name(comp.name) == norm:
return comp
if len(name) >= 12:
prefix = name[:20].lower()
for comp in Component.query.filter(Component.is_deleted.is_(False)).all():
if prefix in (comp.name or "").lower():
return comp
return None
def _enrich_wesp_component(component: Component, row: dict[str, Any]) -> None:
"""Copy zootech nutrients from tab; WESP id/name/dry_matter/price stay canonical."""
from app.lab.services.component_nutrients import save_component_nutrients
raw = row.get("nutrients")
if isinstance(raw, str):
try:
raw = json.loads(raw)
except json.JSONDecodeError:
raw = {}
if not isinstance(raw, dict):
raw = {}
save_component_nutrients(component.id, raw, user_id="reference-db-import")
if row.get("external_no") is not None:
component.external_no = row.get("external_no")
component.updated_by = "tab-enrich"
def _purge_tab_imported_components(stats: ImportStats) -> None:
"""Remove components created from tab feed_ingredients (not used in WESP calc)."""
tab_ids = [
c.id
for c in Component.query.filter(
Component.created_by == "tab-import", Component.is_deleted.is_(False)
).all()
]
if not tab_ids:
return
for comp in Component.query.filter(Component.id.in_(tab_ids)).all():
comp.soft_delete("tab-import-cleanup")
stats.tab_components_purged += 1
def _import_feed_ingredients(conn, stats: ImportStats) -> None:
"""Map tab feed_ingredient → WESP component; enrich nutrients only."""
rows = _fetch_all(
conn,
"""
SELECT id, external_no, name, price_per_kg, dry_matter, nutrients, row_index
FROM feed_ingredients
ORDER BY row_index NULLS LAST, external_no NULLS LAST
""",
)
for row in rows:
component = _find_wesp_component(row)
if component is None:
stats.components_unmatched += 1
continue
_enrich_wesp_component(component, row)
stats.components_enriched += 1
def build_feed_ingredient_mapping(rows: list[dict[str, Any]], stats: ImportStats | None = None) -> dict[str, str | None]:
st = stats or ImportStats()
mapping: dict[str, str | None] = {}
for row in rows:
tab_id = str(row.get("id") or "")
component = _find_wesp_component(row)
if component is None:
st.components_unmatched += 1
mapping[tab_id] = None
continue
_enrich_wesp_component(component, row)
mapping[tab_id] = component.id
st.components_enriched += 1
return mapping
def apply_animal_profile_rows(rows: list[dict[str, Any]], stats: ImportStats | None = None) -> ImportStats:
"""Импорт профилей из списка строк (PG-формат или fixtures JSON)."""
st = stats or ImportStats()
for row in rows:
_upsert_animal_profile_row(row, st)
return st
def _upsert_animal_profile_row(row: dict[str, Any], stats: ImportStats) -> None:
profile = LabAnimalProfile.query.get(row["id"])
if profile is None:
profile = LabAnimalProfile(id=row["id"])
db.session.add(profile)
profile.profile_key = row["key"]
profile.label = row["label"]
profile.ration_type = str(row["type"])
norms_raw = row.get("norms_data")
if isinstance(norms_raw, dict):
save_norms_from_payload(profile, norms_raw)
else:
import_norms_from_legacy_text(profile, norms_raw)
profile.created_at = _dt(row.get("created_at")) or profile.created_at
profile.updated_at = _dt(row.get("updated_at")) or profile.updated_at
profile.created_by = "tab-import"
profile.updated_by = "tab-import"
stats.animal_profiles += 1
def apply_feed_ingredient_rows(rows: list[dict[str, Any]], stats: ImportStats | None = None) -> ImportStats:
"""Обогащение WESP component из строк feed_ingredients (без PG)."""
st = stats or ImportStats()
for row in rows:
component = _find_wesp_component(row)
if component is None:
st.components_unmatched += 1
continue
_enrich_wesp_component(component, row)
st.components_enriched += 1
return st
def _import_animal_profiles(conn, stats: ImportStats) -> None:
rows = _fetch_all(
conn,
"""
SELECT id, key, label, type, norms_data, created_at, updated_at
FROM animal_profiles
ORDER BY key
""",
)
apply_animal_profile_rows(rows, stats)
def _recipe_by_name(name: str) -> Recipe | None:
if not name:
return None
return Recipe.query.filter(Recipe.name == name, Recipe.is_deleted.is_(False)).first()
def _link_recipes_from_tab_projects(conn, stats: ImportStats) -> None:
"""Только ration_type на recipe по имени — без staging-таблиц."""
rows = _fetch_all(
conn,
"""
SELECT name, type
FROM ration_projects
ORDER BY updated_at DESC NULLS LAST
""",
)
seen: set[str] = set()
for row in rows:
name = (row.get("name") or "").strip()
if not name or name in seen:
continue
linked = _recipe_by_name(name)
if linked is None:
continue
if row.get("type"):
linked.ration_type = str(row["type"])
seen.add(name)
stats.recipe_links += 1
def import_from_reference_db(pg_url: str | None = None) -> ImportStats:
"""Full dump from tab PostgreSQL into WESP SQLite."""
stats = ImportStats()
url = pg_url or _pg_url()
conn = psycopg2.connect(url)
db.session.info[WESP_SUPPRESS_SYNC_ENQUEUE] = True
try:
_purge_tab_imported_components(stats)
db.session.flush()
_import_feed_ingredients(conn, stats)
db.session.flush()
_import_animal_profiles(conn, stats)
_link_recipes_from_tab_projects(conn, stats)
db.session.commit()
except Exception as exc:
db.session.rollback()
stats.errors.append(str(exc))
raise
finally:
db.session.info.pop(WESP_SUPPRESS_SYNC_ENQUEUE, None)
conn.close()
return stats