update sync data from site dashboard and testing closing cycles

This commit is contained in:
Alberto-Audrix committed 2026-09-15 16:52:41 +07:00
1 parent 36639d8c38
commit f49e72592c
42 files changed
+1966 -688

No files matched your search

+1
View File
@@ -0,0 +1 @@
"""Active-site mirror pull services."""
@@ -0,0 +1,198 @@
"""HTTP client for farm-location pusat export APIs."""
from __future__ import annotations
from typing import Any, Callable, Iterator
from urllib.parse import urljoin
import httpx
from apps.sync.models import ActiveSite
DEFAULT_PAGE_SIZE = 500
DEFAULT_TIMEOUT = 60.0
OPS_SERIES_KEYS = (
"kpis",
"manual_inputs",
"feed_sacks",
"chicken_countings",
"chicken_weights",
"ai_insights",
)
class FarmExportError(Exception):
def __init__(self, message: str, *, status_code: int | None = None):
self.message = message
self.status_code = status_code
super().__init__(message)
class FarmExportClient:
def __init__(
self,
active_site: ActiveSite,
*,
timeout: float = DEFAULT_TIMEOUT,
transport: httpx.BaseTransport | None = None,
):
base = (active_site.api_base_url or "").strip().rstrip("/") + "/"
if not base or base == "/":
raise FarmExportError("ActiveSite.api_base_url is required.")
if not (active_site.api_key or "").strip():
raise FarmExportError("ActiveSite.api_key is required.")
self._base = base
self._headers = {
"X-API-Key": active_site.api_key.strip(),
"Accept": "application/json",
}
self._timeout = timeout
self._transport = transport
def _url(self, path: str) -> str:
return urljoin(self._base, path.lstrip("/"))
def _get(self, path: str, params: dict[str, Any] | None = None) -> dict[str, Any]:
cleaned = {
key: value
for key, value in (params or {}).items()
if value is not None and value != ""
}
try:
with httpx.Client(
timeout=self._timeout, transport=self._transport
) as client:
response = client.get(
self._url(path), headers=self._headers, params=cleaned
)
except httpx.HTTPError as exc:
raise FarmExportError(f"Network error calling farm export: {exc}") from exc
if response.status_code == 401:
raise FarmExportError("Farm rejected API key (401).", status_code=401)
if response.status_code == 403:
raise FarmExportError("Farm denied access (403).", status_code=403)
if response.status_code >= 400:
detail = (response.text or "")[:300]
raise FarmExportError(
f"Farm export failed ({response.status_code}): {detail}",
status_code=response.status_code,
)
try:
payload = response.json()
except ValueError as exc:
raise FarmExportError("Farm export returned non-JSON body.") from exc
if not isinstance(payload, dict):
raise FarmExportError("Farm export JSON must be an object.")
return payload
def fetch_ops_export(self, **params: Any) -> dict[str, Any]:
return self._get("pusat/export/", params)
def fetch_iot_export(self, **params: Any) -> dict[str, Any]:
return self._get("pusat/export/iot/", params)
def iter_paginated_list(
self,
fetch: Callable[..., dict[str, Any]],
*,
list_key: str,
page_size: int = DEFAULT_PAGE_SIZE,
**params: Any,
) -> Iterator[list[dict[str, Any]]]:
"""Yield result pages until page * page_size >= count for list_key."""
page = 1
while True:
payload = fetch(page=page, page_size=page_size, **params)
data = payload.get("data") or {}
block = data.get(list_key) or {}
results = block.get("results") or []
if not isinstance(results, list):
raise FarmExportError(f"data.{list_key}.results must be a list.")
yield [row for row in results if isinstance(row, dict)]
count = int(block.get("count") or 0)
size = int(block.get("page_size") or page_size)
current = int(block.get("page") or page)
if size <= 0:
break
if current * size >= count or not results:
break
page = current + 1
def collect_ops_series(
self,
*,
site_id: str,
cycle_id: int | None = None,
date_from: str | None = None,
date_to: str | None = None,
page_size: int = DEFAULT_PAGE_SIZE,
) -> dict[str, list[dict[str, Any]]]:
"""Fetch all pages for every ops series key (shared page cursor)."""
collected = {key: [] for key in OPS_SERIES_KEYS}
page = 1
while True:
params: dict[str, Any] = {
"site_id": site_id,
"page": page,
"page_size": page_size,
}
if cycle_id is not None:
params["cycle_id"] = cycle_id
if date_from is not None:
params["from"] = date_from
if date_to is not None:
params["to"] = date_to
payload = self.fetch_ops_export(**params)
data = payload.get("data") or {}
any_more = False
for key in OPS_SERIES_KEYS:
block = data.get(key) or {}
results = block.get("results") or []
if isinstance(results, list):
collected[key].extend(
row for row in results if isinstance(row, dict)
)
count = int(block.get("count") or 0)
size = int(block.get("page_size") or page_size)
current = int(block.get("page") or page)
if size > 0 and current * size < count and results:
any_more = True
if not any_more:
break
page += 1
return collected
def collect_iot_export(
self,
*,
site_id: str,
date_from: str,
date_to: str,
page_size: int = DEFAULT_PAGE_SIZE,
flock_id: int | None = None,
kandang_id: int | None = None,
) -> tuple[dict[str, Any], list[dict[str, Any]]]:
"""Return structure envelope (page 1) plus all iot_panels rows."""
params: dict[str, Any] = {
"site_id": site_id,
"from": date_from,
"to": date_to,
}
if flock_id is not None:
params["flock_id"] = flock_id
if kandang_id is not None:
params["kandang_id"] = kandang_id
first = self.fetch_iot_export(page=1, page_size=page_size, **params)
panels: list[dict[str, Any]] = []
for page_rows in self.iter_paginated_list(
self.fetch_iot_export,
list_key="iot_panels",
page_size=page_size,
**params,
):
panels.extend(page_rows)
return first, panels
+733
View File
@@ -0,0 +1,733 @@
"""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:
count = 0
for raw in rows:
row_date = _as_date(raw.get("date"))
section = str(raw.get("section") or "")[:100]
session = str(raw.get("session") or "")[:100]
if row_date is None:
continue
AIInsight.objects.update_or_create(
cycle=cycle,
date=row_date,
section=section,
session=session,
defaults={
"insight_text": str(raw.get("insight_text") or ""),
"alert": str(raw.get("alert") or "")[:100],
},
)
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
@@ -0,0 +1,123 @@
import httpx
from django.test import SimpleTestCase, TestCase
from apps.accounts.models import User
from apps.sync.models import ActiveSite
from apps.sync.services.farm_export_client import FarmExportClient, FarmExportError
def _ops_page(*, page: int, page_size: int, count: int, results: list) -> dict:
return {
"sites": [],
"data": {
"kpis": {
"count": count,
"page": page,
"page_size": page_size,
"results": results,
},
"manual_inputs": {"count": 0, "page": page, "page_size": page_size, "results": []},
"feed_sacks": {"count": 0, "page": page, "page_size": page_size, "results": []},
"chicken_countings": {
"count": 0,
"page": page,
"page_size": page_size,
"results": [],
},
"chicken_weights": {
"count": 0,
"page": page,
"page_size": page_size,
"results": [],
},
"ai_insights": {"count": 0, "page": page, "page_size": page_size, "results": []},
},
}
class FarmExportClientTests(TestCase):
def setUp(self):
self.gm = User.objects.create_user(
user_name="gm-export", password="x", status=User.STATUS_GM
)
self.active = ActiveSite.objects.create(
code="sukawarna",
name="Sukawarna",
source_site_id="sukawarna-0002",
api_base_url="https://farm.example/api/v1",
api_key="site-secret",
managed_by=self.gm,
)
def test_401_raises_farm_export_error(self):
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(401, json={"detail": "nope"})
client = FarmExportClient(
self.active, transport=httpx.MockTransport(handler)
)
with self.assertRaises(FarmExportError) as ctx:
client.fetch_ops_export(site_id="sukawarna-0002")
self.assertEqual(ctx.exception.status_code, 401)
def test_collect_ops_series_paginates_until_done(self):
calls = []
def handler(request: httpx.Request) -> httpx.Response:
page = int(request.url.params.get("page", "1"))
calls.append(page)
self.assertEqual(request.url.params.get("site_id"), "sukawarna-0002")
self.assertEqual(request.url.params.get("cycle_id"), "9")
self.assertEqual(request.headers.get("X-API-Key"), "site-secret")
if page == 1:
body = _ops_page(
page=1,
page_size=2,
count=3,
results=[{"date": "2026-09-01"}, {"date": "2026-09-02"}],
)
else:
body = _ops_page(
page=2,
page_size=2,
count=3,
results=[{"date": "2026-09-03"}],
)
return httpx.Response(200, json=body)
client = FarmExportClient(
self.active, transport=httpx.MockTransport(handler)
)
series = client.collect_ops_series(
site_id="sukawarna-0002", cycle_id=9, page_size=2
)
self.assertEqual(calls, [1, 2])
self.assertEqual(len(series["kpis"]), 3)
def test_missing_api_key_rejected(self):
self.active.api_key = ""
with self.assertRaises(FarmExportError):
FarmExportClient(self.active)
class SelectCycleIdsUnitTests(SimpleTestCase):
def test_selects_open_and_last_three_closed(self):
from apps.sync.services.mirror_pull import select_source_cycle_ids
rows = [
{"id": 1, "kandang": 10, "status": "active"},
{"id": 2, "kandang": 10, "status": "pending_close"},
{"id": 3, "kandang": 10, "status": "closed", "end_date": "2026-01-01"},
{"id": 4, "kandang": 10, "status": "closed", "end_date": "2026-02-01"},
{"id": 5, "kandang": 10, "status": "closed", "end_date": "2026-03-01"},
{"id": 6, "kandang": 10, "status": "closed", "end_date": "2026-04-01"},
{"id": 7, "kandang": 11, "status": "closed", "end_date": "2026-05-01"},
]
selected = select_source_cycle_ids(rows, closed_limit=3)
self.assertIn(1, selected)
self.assertIn(2, selected)
self.assertIn(6, selected)
self.assertIn(5, selected)
self.assertIn(4, selected)
self.assertNotIn(3, selected)
self.assertIn(7, selected)
+398
View File
@@ -0,0 +1,398 @@
import httpx
from django.test import TestCase
from rest_framework.test import APIClient
from apps.accounts.models import User
from apps.farms.models import Cycle, Flock, Kandang, Site
from apps.operations.models import AIInsight, FeedSacks, IotPanel, KPI
from apps.sync.models import ActiveSite
from apps.sync.services.farm_export_client import FarmExportClient
from apps.sync.services.mirror_pull import sync_active_site_safe
from datetime import date
from unittest.mock import patch
def _empty_series(page=1, page_size=500):
block = {"count": 0, "page": page, "page_size": page_size, "results": []}
return {
"kpis": dict(block),
"manual_inputs": dict(block),
"feed_sacks": dict(block),
"chicken_countings": dict(block),
"chicken_weights": dict(block),
"ai_insights": dict(block),
}
class MirrorPullTests(TestCase):
def setUp(self):
self.gm = User.objects.create_user(
user_name="gm-mirror", password="x", status=User.STATUS_GM
)
self.active = ActiveSite.objects.create(
code="sukawarna",
name="Sukawarna",
source_site_id="sukawarna-0002",
api_base_url="https://farm.example/api/v1",
api_key="site-secret",
managed_by=self.gm,
is_active=True,
)
self.start = date(2026, 9, 1)
def _handler(self):
def handler(request: httpx.Request) -> httpx.Response:
path = request.url.path
params = request.url.params
if path.endswith("/pusat/export/iot/"):
body = {
"sites": [
{
"site_id": "sukawarna-0002",
"site_name": "Sukawarna",
"kandangs": [
{
"kandang_id": 101,
"kandang_name": "K1",
"flocks": [
{"flock_id": 201, "flock_name": "Lantai 1"}
],
}
],
}
],
"data": {
"iot_panels": {
"count": 1,
"page": 1,
"page_size": 500,
"results": [
{
"flock_id": 201,
"date": "2026-09-02",
"timestamp": "2026-09-02T10:00:00+07:00",
"wind_speed": 1.2,
"humidity": 70.0,
"water_total": 10.0,
"average_temperature": 28.0,
"experience_temperature": 27.0,
}
],
}
},
}
return httpx.Response(200, json=body)
cycle_id = params.get("cycle_id")
if cycle_id:
series = _empty_series()
series["feed_sacks"] = {
"count": 1,
"page": 1,
"page_size": 500,
"results": [
{
"date": "2026-09-02",
"in_today": 5,
"out_today": 0,
"in_total": 5,
"out_total": 0,
"feed_use_today": 2,
"feed_use_total": 2,
"cycle": int(cycle_id),
}
],
}
series["chicken_countings"] = {
"count": 1,
"page": 1,
"page_size": 500,
"results": [
{
"date": "2026-09-02",
"total_count": 9900,
"mortality_count": 10,
"cycle": int(cycle_id),
}
],
}
series["chicken_weights"] = {
"count": 1,
"page": 1,
"page_size": 500,
"results": [
{
"date": "2026-09-02",
"age": 1,
"doc_weight": 42,
"average_weight": 50.0,
"chicken_count": 100,
"uniformity": 90.0,
"average_daily_gain": 8.0,
"cycle": int(cycle_id),
}
],
}
series["manual_inputs"] = {
"count": 1,
"page": 1,
"page_size": 500,
"results": [
{
"date": "2026-09-02",
"age_manual": 1,
"mortality_manual": 0,
"cycle": int(cycle_id),
}
],
}
series["ai_insights"] = {
"count": 1,
"page": 1,
"page_size": 500,
"results": [
{
"date": "2026-09-02",
"section": "health",
"session": "am",
"insight_text": "ok",
"alert": "none",
"cycle": int(cycle_id),
}
],
}
series["kpis"] = {
"count": 1,
"page": 1,
"page_size": 500,
"results": [
{
"date": "2026-09-02",
"age": 1,
"fcr": 1.2,
"eef": 300.0,
"chicken_life": 9900,
"cycle": int(cycle_id),
}
],
}
body = {
"site_id": "sukawarna-0002",
"cycle_id": int(cycle_id),
"sites": [],
"data": series,
}
return httpx.Response(200, json=body)
# Structure-only call
body = {
"site_id": "sukawarna-0002",
"sites": [
{
"site_id": "sukawarna-0002",
"site_name": "Sukawarna",
"kandangs": [
{
"kandang_id": 101,
"kandang_name": "K1",
"feed_in_button_urls": [],
"cycles": [
{
"id": 55,
"kandang": 101,
"status": "active",
"start_date": self.start.isoformat(),
"end_date": None,
"chick_in_weight": 42,
"doc_in_count": 10000,
"total_days": 1,
"feed_initial_balance": 120,
"feed_initial_balance_date": self.start.isoformat(),
"feed_initial_balance_iot": 115,
"feed_initial_balance_accuracy": 95.8,
"feed_initial_balance_sync_error": None,
},
{
"id": 50,
"kandang": 101,
"status": "closed",
"start_date": "2026-01-01",
"end_date": "2026-02-01",
"chick_in_weight": 40,
"doc_in_count": 9000,
"total_days": 32,
},
{
"id": 49,
"kandang": 101,
"status": "closed",
"start_date": "2025-10-01",
"end_date": "2025-11-01",
"chick_in_weight": 40,
"doc_in_count": 9000,
"total_days": 32,
},
{
"id": 48,
"kandang": 101,
"status": "closed",
"start_date": "2025-07-01",
"end_date": "2025-08-01",
"chick_in_weight": 40,
"doc_in_count": 9000,
"total_days": 32,
},
{
"id": 47,
"kandang": 101,
"status": "closed",
"start_date": "2025-04-01",
"end_date": "2025-05-01",
"chick_in_weight": 40,
"doc_in_count": 9000,
"total_days": 32,
},
],
}
],
}
],
"data": _empty_series(),
}
return httpx.Response(200, json=body)
return handler
def test_mirror_upsert_idempotent_and_selects_closed_limit(self):
transport = httpx.MockTransport(self._handler())
client = FarmExportClient(self.active, transport=transport)
summary = sync_active_site_safe(
self.active, client=client, closed_limit=3, page_size=500
)
self.assertEqual(summary.status, "ok", summary.message)
self.assertEqual(Cycle.objects.filter(active_site=self.active).count(), 5)
# active + last 3 closed = 4 selected (47 dropped)
self.assertEqual(summary.selected_cycles, 4)
self.assertFalse(
Cycle.objects.filter(active_site=self.active, source_cycle_id=47).exists()
and Cycle.objects.get(
active_site=self.active, source_cycle_id=47
).feed_sacks.exists()
)
# Oldest closed is still created from structure but not series-synced
self.assertTrue(
Cycle.objects.filter(active_site=self.active, source_cycle_id=47).exists()
)
active_cycle = Cycle.objects.get(active_site=self.active, source_cycle_id=55)
self.assertEqual(active_cycle.feed_initial_balance, 120)
self.assertEqual(active_cycle.feed_initial_balance_date, self.start)
self.assertEqual(active_cycle.feed_initial_balance_iot, 115)
self.assertEqual(active_cycle.feed_initial_balance_accuracy, 95.8)
self.assertIsNone(active_cycle.feed_initial_balance_sync_error)
self.assertEqual(FeedSacks.objects.filter(cycle=active_cycle).count(), 1)
self.assertEqual(KPI.objects.filter(cycle=active_cycle).count(), 1)
self.assertEqual(AIInsight.objects.filter(cycle=active_cycle).count(), 1)
self.assertEqual(Flock.objects.filter(source_flock_id=201).count(), 1)
self.assertEqual(IotPanel.objects.count(), 1)
# Second sync is idempotent
summary2 = sync_active_site_safe(
self.active, client=FarmExportClient(self.active, transport=transport)
)
self.assertEqual(summary2.status, "ok", summary2.message)
self.assertEqual(FeedSacks.objects.filter(cycle=active_cycle).count(), 1)
self.assertEqual(IotPanel.objects.count(), 1)
self.active.refresh_from_db()
self.assertEqual(self.active.last_sync_status, "ok")
def test_sync_api_endpoint(self):
transport = httpx.MockTransport(self._handler())
api = APIClient()
api.force_authenticate(self.gm)
with patch(
"apps.sync.services.mirror_pull.FarmExportClient",
side_effect=lambda active_site, **kwargs: FarmExportClient(
active_site, transport=transport
),
):
response = api.post(f"/api/v1/active-sites/{self.active.pk}/sync/")
self.assertEqual(response.status_code, 200, response.data)
self.assertEqual(response.data["summary"]["status"], "ok")
self.active.refresh_from_db()
self.assertEqual(self.active.last_sync_status, "ok")
class CloseIngestSourceKandangTests(TestCase):
def setUp(self):
self.gm = User.objects.create_user(
user_name="gm-ingest", password="x", status=User.STATUS_GM
)
self.machine = User.objects.create_user(
user_name="bot", password="x", status=User.STATUS_GM
)
self.active = ActiveSite.objects.create(
code="site-a",
name="Site A",
source_site_id="site-a-0001",
api_base_url="https://a.example/api/v1",
api_key="k",
managed_by=self.gm,
)
self.client = APIClient()
self.client.force_authenticate(self.machine)
def test_ingest_prefers_source_kandang_id(self):
response = self.client.post(
"/api/v1/cycle-close/ingest/",
{
"active_site_id": self.active.pk,
"source_cycle_id": 99,
"source_kandang_id": 777,
"kandang_name": "Kandang Baru",
"site_name": "Site A",
"start_date": "2026-09-01",
"proposed_end_date": "2026-09-15",
"doc_in_weight": 40,
"doc_in_count": 1000,
},
format="json",
)
self.assertEqual(response.status_code, 200, response.data)
kandang = Kandang.objects.get(source_kandang_id=777)
self.assertEqual(kandang.kandang_name, "Kandang Baru")
cycle = Cycle.objects.get(active_site=self.active, source_cycle_id=99)
self.assertEqual(cycle.kandang_id, kandang.pk)
class AIInsightUniqueTests(TestCase):
def setUp(self):
owner = User.objects.create_user(user_name="owner", password="x")
site = Site.objects.create(site_id="s-1", site_name="S", user=owner)
kandang = Kandang.objects.create(site=site, kandang_name="K")
self.cycle = Cycle.objects.create(
kandang=kandang,
start_date=date.today(),
doc_in_weight=40,
doc_in_count=100,
)
def test_unique_constraint_allows_upsert_key(self):
AIInsight.objects.create(
cycle=self.cycle,
date=date.today(),
section="a",
session="am",
insight_text="one",
alert="",
)
AIInsight.objects.update_or_create(
cycle=self.cycle,
date=date.today(),
section="a",
session="am",
defaults={"insight_text": "two", "alert": "x"},
)
self.assertEqual(AIInsight.objects.filter(cycle=self.cycle).count(), 1)
self.assertEqual(
AIInsight.objects.get(cycle=self.cycle).insight_text, "two"
)
+31 -8
View File
@@ -19,7 +19,7 @@ class ActiveSiteViewSet(viewsets.ModelViewSet):
"""
List/read: scoped (director all; BUH managed GMs; GM own).
Create (claim): GM and GM Admin claim a site-registered ActiveSite by code.
Sync: GM and GM Admin only (stubs until mirror pull is implemented).
Sync: GM and GM Admin only (pulls farm pusat export into HQ mirror).
Site dashboards create/fill ActiveSite via register-from-site (API key auth).
Claiming attaches managed_by and ensures an HQ Site for the view list.
@@ -203,24 +203,47 @@ class ActiveSiteViewSet(viewsets.ModelViewSet):
{"detail": "Active site is inactive."},
status=status.HTTP_400_BAD_REQUEST,
)
if not active_site.api_base_url or not active_site.api_key:
return Response(
{"detail": "Active site is missing api_base_url or api_key."},
status=status.HTTP_400_BAD_REQUEST,
)
if not (active_site.source_site_id or "").strip():
return Response(
{"detail": "Active site is missing source_site_id."},
status=status.HTTP_400_BAD_REQUEST,
)
from apps.sync.services.mirror_pull import STATUS_OK, sync_active_site_safe
summary = sync_active_site_safe(active_site)
active_site.refresh_from_db()
http_status = (
status.HTTP_200_OK
if summary.status == STATUS_OK
else status.HTTP_502_BAD_GATEWAY
)
return Response(
{
"detail": "Active site sync not implemented yet.",
"detail": summary.message or "Sync finished.",
"summary": summary.to_dict(),
"active_site": ActiveSiteSerializer(active_site).data,
},
status=status.HTTP_501_NOT_IMPLEMENTED,
status=http_status,
)
@action(detail=False, methods=["post"], url_path="sync-all")
def sync_all(self, request):
if not can_sync(request.user):
return Response({"detail": "Sync not allowed."}, status=status.HTTP_403_FORBIDDEN)
from apps.sync.services.mirror_pull import STATUS_OK, sync_all_active_sites
sites = list(self.get_queryset().filter(is_active=True))
summaries = sync_all_active_sites(queryset=sites)
ok = sum(1 for s in summaries if s.status == STATUS_OK)
return Response(
{
"detail": "Active site sync-all not implemented yet.",
"active_site_count": len(sites),
"active_sites": ActiveSiteSerializer(sites, many=True).data,
},
status=status.HTTP_501_NOT_IMPLEMENTED,
"detail": f"Synced {ok}/{len(summaries)} active site(s).",
"active_site_count": len(summaries),
"summaries": [s.to_dict() for s in summaries],
}
)