Files
dashboard-cpsp-executive/backend/apps/sync/services/mirror_pull.py
T

749 lines
26 KiB
Python

"""Pull farm ops + IoT exports and upsert into the HQ mirror."""
from __future__ import annotations
import logging
from dataclasses import dataclass, field
from datetime import date, datetime
from typing import Any
from django.db import transaction
from django.utils import timezone
from django.utils.dateparse import parse_date, parse_datetime
from apps.farms.models import Cycle, Flock, Kandang, Site
from apps.farms.services.active_site_registry import ensure_hq_site_for_active_site
from apps.operations.models import (
AIInsight,
ChickenCounting,
ChickenWeight,
FeedSacks,
IotPanel,
KPI,
ManualInput,
)
from apps.sync.models import ActiveSite
from apps.sync.services.farm_export_client import (
DEFAULT_PAGE_SIZE,
FarmExportClient,
FarmExportError,
)
logger = logging.getLogger(__name__)
CLOSED_LIMIT_PER_KANDANG = 3
STATUS_OK = "ok"
STATUS_ERROR = "error"
@dataclass
class MirrorSyncSummary:
active_site_id: int
status: str = STATUS_OK
message: str = ""
kandangs: int = 0
cycles: int = 0
flocks: int = 0
selected_cycles: int = 0
series_rows: dict[str, int] = field(default_factory=dict)
iot_panels: int = 0
def to_dict(self) -> dict[str, Any]:
return {
"active_site_id": self.active_site_id,
"status": self.status,
"message": self.message,
"kandangs": self.kandangs,
"cycles": self.cycles,
"flocks": self.flocks,
"selected_cycles": self.selected_cycles,
"series_rows": self.series_rows,
"iot_panels": self.iot_panels,
}
def _as_int(value: Any) -> int | None:
if value is None or value == "":
return None
try:
return int(value)
except (TypeError, ValueError):
return None
def _as_float(value: Any) -> float | None:
if value is None or value == "":
return None
try:
return float(value)
except (TypeError, ValueError):
return None
def _as_date(value: Any) -> date | None:
if value is None or value == "":
return None
if isinstance(value, date) and not isinstance(value, datetime):
return value
if isinstance(value, datetime):
return value.date()
return parse_date(str(value)[:10])
def _as_datetime(value: Any):
if value is None or value == "":
return None
if isinstance(value, datetime):
return value
return parse_datetime(str(value))
def _float(value: Any, default: float = 0.0) -> float:
try:
return float(value)
except (TypeError, ValueError):
return default
def _int(value: Any, default: int = 0) -> int:
try:
return int(value)
except (TypeError, ValueError):
return default
def select_source_cycle_ids(
structure_cycles: list[dict[str, Any]],
*,
closed_limit: int = CLOSED_LIMIT_PER_KANDANG,
) -> list[int]:
"""Pick farm cycle ids: all active/pending_close + last N closed per kandang."""
by_kandang: dict[int, list[dict[str, Any]]] = {}
selected: list[int] = []
for raw in structure_cycles:
source_id = _as_int(raw.get("id") or raw.get("cycle_id"))
if source_id is None:
continue
status = str(raw.get("status") or "").strip().lower()
kandang_id = _as_int(raw.get("kandang")) or 0
if status in (Cycle.STATUS_ACTIVE, Cycle.STATUS_PENDING_CLOSE):
selected.append(source_id)
continue
if status == Cycle.STATUS_CLOSED:
by_kandang.setdefault(kandang_id, []).append(raw)
for rows in by_kandang.values():
rows.sort(
key=lambda r: (_as_date(r.get("end_date")) or date.min, _as_int(r.get("id")) or 0),
reverse=True,
)
for raw in rows[:closed_limit]:
source_id = _as_int(raw.get("id") or raw.get("cycle_id"))
if source_id is not None:
selected.append(source_id)
# Stable unique order
seen: set[int] = set()
ordered: list[int] = []
for cid in selected:
if cid not in seen:
seen.add(cid)
ordered.append(cid)
return ordered
def _upsert_kandang(site: Site, raw: dict[str, Any]) -> Kandang:
source_id = _as_int(raw.get("kandang_id") or raw.get("id"))
name = str(raw.get("kandang_name") or "Kandang")[:30]
feed_urls = raw.get("feed_in_button_urls")
if not isinstance(feed_urls, list):
feed_urls = []
if source_id is not None:
kandang = Kandang.objects.filter(site=site, source_kandang_id=source_id).first()
if kandang is None:
kandang = Kandang(
site=site,
source_kandang_id=source_id,
kandang_name=name,
feed_in_button_urls=feed_urls,
)
kandang.save()
else:
kandang.kandang_name = name
kandang.feed_in_button_urls = feed_urls
kandang.save(
update_fields=["kandang_name", "feed_in_button_urls", "updated_at"]
)
return kandang
kandang, _ = Kandang.objects.get_or_create(
site=site, kandang_name=name, defaults={"feed_in_button_urls": feed_urls}
)
return kandang
def _upsert_cycle(
active_site: ActiveSite,
kandang: Kandang,
raw: dict[str, Any],
) -> Cycle | None:
source_id = _as_int(raw.get("id") or raw.get("cycle_id"))
start = _as_date(raw.get("start_date"))
if source_id is None or start is None:
return None
doc_in_weight = _int(raw.get("chick_in_weight", raw.get("doc_in_weight")), 0)
doc_in_count = _int(raw.get("doc_in_count"), 0)
end = _as_date(raw.get("end_date"))
status = str(raw.get("status") or Cycle.STATUS_ACTIVE).strip().lower()
if status not in {
Cycle.STATUS_ACTIVE,
Cycle.STATUS_PENDING_CLOSE,
Cycle.STATUS_CLOSED,
}:
status = Cycle.STATUS_ACTIVE
defaults = {
"kandang": kandang,
"start_date": start,
"end_date": end,
"doc_in_weight": doc_in_weight,
"doc_in_count": doc_in_count,
"total_days": max(_int(raw.get("total_days"), 1), 1),
"feed_initial_balance": _int(raw.get("feed_initial_balance"), 0),
"feed_initial_balance_date": _as_date(raw.get("feed_initial_balance_date")),
"feed_initial_balance_iot": _as_int(raw.get("feed_initial_balance_iot")),
"feed_initial_balance_accuracy": _as_float(
raw.get("feed_initial_balance_accuracy")
),
"feed_initial_balance_sync_error": (
str(raw.get("feed_initial_balance_sync_error"))
if raw.get("feed_initial_balance_sync_error") not in (None, "")
else None
),
"status": status,
"proposed_end_date": _as_date(raw.get("proposed_end_date")),
}
cycle = Cycle.objects.filter(
active_site=active_site, source_cycle_id=source_id
).first()
if cycle is None:
cycle = Cycle(active_site=active_site, source_cycle_id=source_id, **defaults)
cycle.save()
else:
for key, value in defaults.items():
setattr(cycle, key, value)
cycle.save()
return cycle
def _upsert_flock(kandang: Kandang, raw: dict[str, Any]) -> Flock | None:
source_id = _as_int(raw.get("flock_id") or raw.get("id"))
name = str(raw.get("flock_name") or "Flock")[:30]
if source_id is None:
flock, _ = Flock.objects.get_or_create(
kandang=kandang, flock_name=name, defaults={}
)
return flock
flock = Flock.objects.filter(kandang=kandang, source_flock_id=source_id).first()
if flock is None:
flock = Flock(kandang=kandang, source_flock_id=source_id, flock_name=name)
flock.save()
else:
flock.flock_name = name
flock.save(update_fields=["flock_name", "updated_at"])
return flock
def _upsert_feed_sacks(cycle: Cycle, rows: list[dict[str, Any]]) -> int:
count = 0
for raw in rows:
row_date = _as_date(raw.get("date"))
if row_date is None:
continue
FeedSacks.objects.update_or_create(
cycle=cycle,
date=row_date,
defaults={
"in_today": _int(raw.get("in_today")),
"out_today": _int(raw.get("out_today")),
"in_total": _int(raw.get("in_total")),
"out_total": _int(raw.get("out_total")),
"feed_use_today": _int(raw.get("feed_use_today")),
"feed_use_total": _int(raw.get("feed_use_total")),
},
)
count += 1
return count
def _upsert_chicken_countings(cycle: Cycle, rows: list[dict[str, Any]]) -> int:
count = 0
for raw in rows:
row_date = _as_date(raw.get("date"))
if row_date is None:
continue
ChickenCounting.objects.update_or_create(
cycle=cycle,
date=row_date,
defaults={
"total_count": _int(raw.get("total_count")),
"mortality_count": _int(raw.get("mortality_count")),
},
)
count += 1
return count
def _upsert_chicken_weights(cycle: Cycle, rows: list[dict[str, Any]]) -> int:
count = 0
for raw in rows:
row_date = _as_date(raw.get("date"))
if row_date is None:
continue
ChickenWeight.objects.update_or_create(
cycle=cycle,
date=row_date,
defaults={
"age": _int(raw.get("age")),
"doc_weight": _int(raw.get("doc_weight")),
"average_weight": _float(raw.get("average_weight")),
"chicken_count": _int(raw.get("chicken_count")),
"uniformity": _float(raw.get("uniformity")),
"average_daily_gain": _float(raw.get("average_daily_gain")),
},
)
count += 1
return count
def _upsert_manual_inputs(cycle: Cycle, rows: list[dict[str, Any]]) -> int:
count = 0
for raw in rows:
row_date = _as_date(raw.get("date"))
if row_date is None:
continue
ManualInput.objects.update_or_create(
cycle=cycle,
date=row_date,
defaults={
"age_manual": _int(raw.get("age_manual")),
"mortality_manual": _int(raw.get("mortality_manual")),
"mortality_manual_total": _int(raw.get("mortality_manual_total")),
"feed_in_manual": _int(raw.get("feed_in_manual")),
"feed_use_manual": _int(raw.get("feed_use_manual")),
"feed_in_manual_total": _int(raw.get("feed_in_manual_total")),
"feed_use_manual_total": _int(raw.get("feed_use_manual_total")),
"feed_out_manual": _int(raw.get("feed_out_manual")),
"feed_out_manual_total": _int(raw.get("feed_out_manual_total")),
"harvest_manual": _int(raw.get("harvest_manual")),
"harvest_manual_total": _int(raw.get("harvest_manual_total")),
"harvest_weight_manual": _float(raw.get("harvest_weight_manual")),
"harvest_weight_manual_total": _float(
raw.get("harvest_weight_manual_total")
),
"manual_weight": _float(raw.get("manual_weight")),
"average_harvest_day_manual": _float(
raw.get("average_harvest_day_manual")
),
"average_harvest_weight_manual": _float(
raw.get("average_harvest_weight_manual")
),
"fcr_manual": _float(raw.get("fcr_manual")),
"eef_manual": _float(raw.get("eef_manual")),
},
)
count += 1
return count
def _ensure_related_for_kpi(cycle: Cycle, row_date: date) -> tuple:
cc, _ = ChickenCounting.objects.get_or_create(
cycle=cycle, date=row_date, defaults={}
)
cw, _ = ChickenWeight.objects.get_or_create(
cycle=cycle,
date=row_date,
defaults={
"age": 0,
"doc_weight": cycle.doc_in_weight,
"average_weight": 0.0,
"chicken_count": 0,
"uniformity": 0.0,
"average_daily_gain": 0.0,
},
)
fs, _ = FeedSacks.objects.get_or_create(cycle=cycle, date=row_date, defaults={})
manual, _ = ManualInput.objects.get_or_create(
cycle=cycle, date=row_date, defaults={}
)
return cc, cw, fs, manual
def _upsert_kpis(cycle: Cycle, rows: list[dict[str, Any]]) -> int:
count = 0
for raw in rows:
row_date = _as_date(raw.get("date"))
if row_date is None:
continue
cc, cw, fs, manual = _ensure_related_for_kpi(cycle, row_date)
KPI.objects.update_or_create(
cycle=cycle,
date=row_date,
defaults={
"cc": cc,
"cw": cw,
"fs": fs,
"manual": manual,
"age": _int(raw.get("age")),
"mortality": _int(raw.get("mortality")),
"mortality_total": _int(raw.get("mortality_total")),
"feed": _int(raw.get("feed")),
"feed_total": _int(raw.get("feed_total")),
"harvest": _int(raw.get("harvest")),
"harvest_total": _int(raw.get("harvest_total")),
"harvest_weight": _float(raw.get("harvest_weight")),
"harvest_weight_total": _float(raw.get("harvest_weight_total")),
"chicken_life": _int(raw.get("chicken_life")),
"chicken_life_percentage": _float(raw.get("chicken_life_percentage")),
"iot_weight": _float(raw.get("iot_weight")),
"average_harvest_day": _float(raw.get("average_harvest_day")),
"average_harvest_weight": _float(raw.get("average_harvest_weight")),
"fcr": _float(raw.get("fcr")),
"eef": _float(raw.get("eef")),
},
)
count += 1
return count
def _upsert_ai_insights(cycle: Cycle, rows: list[dict[str, Any]]) -> int:
import json
count = 0
for raw in rows:
row_date = _as_date(raw.get("date")) or cycle.start_date
topic = str(raw.get("topic") or "")[:64]
report_type = str(raw.get("report_type") or "page")[:32]
report_period = str(raw.get("report_period") or "current")[:64]
insight_text = str(raw.get("insight_text") or "")
summary = str(raw.get("summary") or "")
alert = str(raw.get("alert") or "")[:100]
raw_citations = raw.get("citations")
if isinstance(raw_citations, (list, dict)):
citations_str = json.dumps(raw_citations, ensure_ascii=False)
elif isinstance(raw_citations, str):
citations_str = raw_citations
else:
citations_str = "[]"
AIInsight.objects.update_or_create(
cycle=cycle,
kandang=cycle.kandang,
topic=topic,
report_type=report_type,
report_period=report_period,
defaults={
"date": row_date,
"insight_text": insight_text,
"summary": summary,
"alert": alert,
"citations": citations_str,
},
)
count += 1
return count
def _upsert_iot_panels(
flock_by_source: dict[int, Flock], rows: list[dict[str, Any]]
) -> int:
count = 0
for raw in rows:
source_flock = _as_int(raw.get("flock_id") or raw.get("flock"))
flock = flock_by_source.get(source_flock) if source_flock is not None else None
if flock is None:
continue
ts = _as_datetime(raw.get("timestamp"))
row_date = _as_date(raw.get("date"))
if row_date is None and ts is not None:
row_date = timezone.localtime(ts).date() if timezone.is_aware(ts) else ts.date()
if row_date is None or ts is None:
continue
IotPanel.objects.update_or_create(
flock=flock,
timestamp=ts,
defaults={
"date": row_date,
"wind_speed": _float(raw.get("wind_speed")),
"humidity": _float(raw.get("humidity")),
"water_total": _float(raw.get("water_total")),
"average_temperature": _float(raw.get("average_temperature")),
"experience_temperature": _float(raw.get("experience_temperature")),
},
)
count += 1
return count
def _mark_active_site(active_site: ActiveSite, summary: MirrorSyncSummary) -> None:
active_site.last_synced_at = timezone.now()
active_site.last_sync_status = summary.status
active_site.last_sync_message = summary.message[:2000]
active_site.save(
update_fields=[
"last_synced_at",
"last_sync_status",
"last_sync_message",
"updated_at",
]
)
@transaction.atomic
def sync_active_site(
active_site: ActiveSite,
*,
client: FarmExportClient | None = None,
closed_limit: int = CLOSED_LIMIT_PER_KANDANG,
page_size: int = DEFAULT_PAGE_SIZE,
) -> MirrorSyncSummary:
summary = MirrorSyncSummary(active_site_id=active_site.pk)
site_id = str(active_site.source_site_id or "").strip()
if not site_id:
summary.status = STATUS_ERROR
summary.message = "ActiveSite.source_site_id is required."
_mark_active_site(active_site, summary)
return summary
try:
site = ensure_hq_site_for_active_site(active_site)
export_client = client or FarmExportClient(active_site)
ops_structure = export_client.fetch_ops_export(
site_id=site_id, page=1, page_size=page_size
)
sites_tree = ops_structure.get("sites") or []
if not sites_tree:
raise FarmExportError("Ops export returned empty sites structure.")
structure_cycles: list[dict[str, Any]] = []
cycle_map: dict[int, Cycle] = {}
kandang_by_source: dict[int, Kandang] = {}
for site_node in sites_tree:
for kandang_raw in site_node.get("kandangs") or []:
if not isinstance(kandang_raw, dict):
continue
kandang = _upsert_kandang(site, kandang_raw)
summary.kandangs += 1
source_kandang = _as_int(
kandang_raw.get("kandang_id") or kandang_raw.get("id")
)
if source_kandang is not None:
kandang_by_source[source_kandang] = kandang
for cycle_raw in kandang_raw.get("cycles") or []:
if not isinstance(cycle_raw, dict):
continue
structure_cycles.append(cycle_raw)
cycle = _upsert_cycle(active_site, kandang, cycle_raw)
if cycle is None:
continue
summary.cycles += 1
if cycle.source_cycle_id is not None:
cycle_map[int(cycle.source_cycle_id)] = cycle
selected_ids = select_source_cycle_ids(
structure_cycles, closed_limit=closed_limit
)
summary.selected_cycles = len(selected_ids)
series_totals = {
"kpis": 0,
"manual_inputs": 0,
"feed_sacks": 0,
"chicken_countings": 0,
"chicken_weights": 0,
"ai_insights": 0,
}
for source_cycle_id in selected_ids:
cycle = cycle_map.get(source_cycle_id)
if cycle is None:
continue
series = export_client.collect_ops_series(
site_id=site_id,
cycle_id=source_cycle_id,
page_size=page_size,
)
# Upsert non-KPI series first so KPI FKs resolve.
series_totals["feed_sacks"] += _upsert_feed_sacks(
cycle, series.get("feed_sacks") or []
)
series_totals["chicken_countings"] += _upsert_chicken_countings(
cycle, series.get("chicken_countings") or []
)
series_totals["chicken_weights"] += _upsert_chicken_weights(
cycle, series.get("chicken_weights") or []
)
series_totals["manual_inputs"] += _upsert_manual_inputs(
cycle, series.get("manual_inputs") or []
)
series_totals["ai_insights"] += _upsert_ai_insights(
cycle, series.get("ai_insights") or []
)
series_totals["kpis"] += _upsert_kpis(cycle, series.get("kpis") or [])
summary.series_rows = series_totals
# IoT window = union of selected cycle spans (explicit from/to).
today = timezone.localdate()
window_from: date | None = None
window_to: date | None = None
for source_cycle_id in selected_ids:
cycle = cycle_map.get(source_cycle_id)
if cycle is None:
continue
start = cycle.start_date
end = cycle.end_date or today
window_from = start if window_from is None else min(window_from, start)
window_to = end if window_to is None else max(window_to, end)
flock_by_source: dict[int, Flock] = {}
if window_from is not None and window_to is not None:
iot_envelope, panels = export_client.collect_iot_export(
site_id=site_id,
date_from=window_from.isoformat(),
date_to=window_to.isoformat(),
page_size=page_size,
)
for site_node in iot_envelope.get("sites") or []:
for kandang_raw in site_node.get("kandangs") or []:
if not isinstance(kandang_raw, dict):
continue
source_kandang = _as_int(
kandang_raw.get("kandang_id") or kandang_raw.get("id")
)
kandang = (
kandang_by_source.get(source_kandang)
if source_kandang is not None
else None
)
if kandang is None:
kandang = _upsert_kandang(site, kandang_raw)
if source_kandang is not None:
kandang_by_source[source_kandang] = kandang
for flock_raw in kandang_raw.get("flocks") or []:
if not isinstance(flock_raw, dict):
continue
flock = _upsert_flock(kandang, flock_raw)
if flock is None:
continue
summary.flocks += 1
source_flock = _as_int(
flock_raw.get("flock_id") or flock_raw.get("id")
)
if source_flock is not None:
flock_by_source[source_flock] = flock
summary.iot_panels = _upsert_iot_panels(flock_by_source, panels)
summary.message = (
f"Synced {summary.selected_cycles} cycle(s), "
f"{sum(series_totals.values())} series row(s), "
f"{summary.iot_panels} IoT panel(s)."
)
_mark_active_site(active_site, summary)
return summary
except Exception as exc:
logger.exception("Mirror sync failed for active_site=%s", active_site.pk)
summary.status = STATUS_ERROR
summary.message = str(getattr(exc, "message", exc))[:2000]
# Avoid swallowing the outer atomic block with a nested failure write —
# mark status in a separate save after rollback via on_commit is awkward;
# write status outside transaction by raising after marking.
raise MirrorSyncFailed(summary) from exc
class MirrorSyncFailed(Exception):
def __init__(self, summary: MirrorSyncSummary):
self.summary = summary
super().__init__(summary.message)
def sync_active_site_safe(
active_site: ActiveSite,
*,
client: FarmExportClient | None = None,
closed_limit: int = CLOSED_LIMIT_PER_KANDANG,
page_size: int = DEFAULT_PAGE_SIZE,
) -> MirrorSyncSummary:
"""Sync one site; capture errors into ActiveSite.last_sync_* without raising."""
try:
return sync_active_site(
active_site,
client=client,
closed_limit=closed_limit,
page_size=page_size,
)
except MirrorSyncFailed as exc:
summary = exc.summary
ActiveSite.objects.filter(pk=active_site.pk).update(
last_synced_at=timezone.now(),
last_sync_status=summary.status,
last_sync_message=summary.message[:2000],
updated_at=timezone.now(),
)
return summary
except Exception as exc:
summary = MirrorSyncSummary(
active_site_id=active_site.pk,
status=STATUS_ERROR,
message=str(exc)[:2000],
)
ActiveSite.objects.filter(pk=active_site.pk).update(
last_synced_at=timezone.now(),
last_sync_status=summary.status,
last_sync_message=summary.message,
updated_at=timezone.now(),
)
return summary
def sync_all_active_sites(
*,
queryset=None,
closed_limit: int = CLOSED_LIMIT_PER_KANDANG,
page_size: int = DEFAULT_PAGE_SIZE,
) -> list[MirrorSyncSummary]:
qs = queryset
if qs is None:
qs = (
ActiveSite.objects.filter(is_active=True)
.exclude(source_site_id="")
.exclude(api_base_url="")
.exclude(api_key="")
)
results: list[MirrorSyncSummary] = []
for active_site in qs:
if not active_site.api_base_url or not active_site.api_key:
summary = MirrorSyncSummary(
active_site_id=active_site.pk,
status=STATUS_ERROR,
message="Missing api_base_url or api_key.",
)
ActiveSite.objects.filter(pk=active_site.pk).update(
last_synced_at=timezone.now(),
last_sync_status=summary.status,
last_sync_message=summary.message,
updated_at=timezone.now(),
)
results.append(summary)
continue
results.append(
sync_active_site_safe(
active_site, closed_limit=closed_limit, page_size=page_size
)
)
return results