132 lines
4.6 KiB
Python
132 lines
4.6 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
from datetime import UTC, datetime
|
|
from typing import Any
|
|
|
|
from sqlalchemy import select
|
|
|
|
from app.core.database import session_scope
|
|
from app.modules.zootech.report_models import ZootechFeedAlert, ZootechLoadingReport, ZootechUnloadingReport
|
|
|
|
REPORT_TABLE_MAP = {
|
|
"loading_report": ZootechLoadingReport,
|
|
"unloading_report": ZootechUnloadingReport,
|
|
"feed_alert": ZootechFeedAlert,
|
|
}
|
|
|
|
|
|
def _write_payload_json(row: Any, extra: dict[str, Any]) -> None:
|
|
if not extra:
|
|
return
|
|
row.payload_json = json.dumps(extra, ensure_ascii=False)
|
|
|
|
|
|
def repair_report_payloads_from_event_log(enterprise_id: str) -> int:
|
|
"""Re-apply latest sync event payload per report row (one-time repair helper)."""
|
|
from app.modules.sync.models import SyncEventLog
|
|
|
|
repaired = 0
|
|
events: list[tuple[str, str, str, int, str, str | None]] = []
|
|
with session_scope() as db:
|
|
rows = list(
|
|
db.scalars(
|
|
select(SyncEventLog)
|
|
.where(
|
|
SyncEventLog.enterprise_id == enterprise_id,
|
|
SyncEventLog.table_name.in_(("loading_report", "unloading_report", "feed_alert")),
|
|
)
|
|
.order_by(SyncEventLog.record_id.asc(), SyncEventLog.received_at.desc())
|
|
)
|
|
)
|
|
latest_by_record: dict[tuple[str, str], SyncEventLog] = {}
|
|
for row in rows:
|
|
key = (row.table_name, row.record_id)
|
|
if key not in latest_by_record:
|
|
latest_by_record[key] = row
|
|
for (table_name, record_id), event in latest_by_record.items():
|
|
events.append(
|
|
(
|
|
table_name,
|
|
record_id,
|
|
event.payload_json or "{}",
|
|
int(event.version or 1),
|
|
str(event.content_hash or ""),
|
|
event.origin_site_id,
|
|
)
|
|
)
|
|
|
|
for table_name, record_id, payload_json, version, content_hash, origin_site_id in events:
|
|
try:
|
|
payload = json.loads(payload_json)
|
|
except json.JSONDecodeError:
|
|
continue
|
|
if not isinstance(payload, dict):
|
|
continue
|
|
apply_report_change(
|
|
enterprise_id,
|
|
table_name,
|
|
record_id,
|
|
"upsert",
|
|
payload,
|
|
version,
|
|
content_hash,
|
|
farm_hub_id=origin_site_id,
|
|
)
|
|
repaired += 1
|
|
return repaired
|
|
|
|
|
|
def apply_report_change(
|
|
enterprise_id: str,
|
|
table_name: str,
|
|
record_id: str,
|
|
action: str,
|
|
payload: dict[str, Any],
|
|
version: int,
|
|
content_hash: str,
|
|
farm_hub_id: str | None = None,
|
|
) -> None:
|
|
model = REPORT_TABLE_MAP.get(table_name)
|
|
if not model:
|
|
return
|
|
with session_scope() as db:
|
|
row = db.scalar(
|
|
select(model).where(model.enterprise_id == enterprise_id, model.id == record_id)
|
|
)
|
|
if action == "delete":
|
|
if row:
|
|
row.is_deleted = True
|
|
row.version = version
|
|
row.content_hash = content_hash
|
|
row.updated_at = datetime.now(UTC)
|
|
return
|
|
data = dict(payload)
|
|
data["id"] = record_id
|
|
data["enterprise_id"] = enterprise_id
|
|
data["farm_hub_id"] = farm_hub_id or data.get("farm_hub_id")
|
|
data["version"] = version
|
|
data["content_hash"] = content_hash
|
|
data["is_deleted"] = False
|
|
if table_name == "feed_alert":
|
|
if data.get("event_type") and not data.get("alert_type"):
|
|
data["alert_type"] = str(data["event_type"])
|
|
if data.get("detail") and not data.get("message"):
|
|
data["message"] = str(data["detail"])
|
|
if row:
|
|
allowed = {c.key for c in model.__table__.columns}
|
|
extra = {k: v for k, v in data.items() if k not in allowed}
|
|
for key, value in data.items():
|
|
if key in allowed and key not in ("enterprise_id", "created_at"):
|
|
setattr(row, key, value)
|
|
if extra and "payload_json" in allowed:
|
|
_write_payload_json(row, extra)
|
|
row.updated_at = datetime.now(UTC)
|
|
else:
|
|
allowed = {c.key for c in model.__table__.columns}
|
|
filtered = {k: v for k, v in data.items() if k in allowed}
|
|
extra = {k: v for k, v in data.items() if k not in allowed}
|
|
if extra and "payload_json" in allowed:
|
|
filtered["payload_json"] = json.dumps(extra, ensure_ascii=False)
|
|
db.add(model(**filtered))
|