215 lines
7.4 KiB
Python
215 lines
7.4 KiB
Python
from __future__ import annotations
|
|
|
|
import re
|
|
from datetime import date
|
|
from typing import Any
|
|
from urllib.parse import urljoin
|
|
|
|
import httpx
|
|
from django.conf import settings
|
|
from django.core.exceptions import ValidationError
|
|
from django.db import transaction
|
|
from django.utils import timezone
|
|
|
|
from apps.farms.models import Cycle
|
|
from apps.operations.models import ChickenCounting
|
|
|
|
|
|
class ChickenCountingEdgeError(Exception):
|
|
def __init__(self, message: str, status_code: int | None = None):
|
|
super().__init__(message)
|
|
self.status_code = status_code
|
|
|
|
|
|
def kandang_index(kandang_name: str) -> str:
|
|
match = re.search(r"(\d+)", kandang_name or "")
|
|
if not match:
|
|
raise ValidationError(
|
|
f"Cannot derive edge kandang index from kandang name: {kandang_name!r}"
|
|
)
|
|
return match.group(1)
|
|
|
|
|
|
def location_prefix_for_kandang(kandang_name: str) -> str:
|
|
"""``Kandang 1`` → ``K1`` (matches floor prefixes like ``K1-L1``)."""
|
|
return f"K{kandang_index(kandang_name)}"
|
|
|
|
|
|
def coop_for_kandang(kandang_name: str) -> str:
|
|
"""``Kandang 1`` → ``K1`` (mortality ``coop`` field)."""
|
|
return location_prefix_for_kandang(kandang_name)
|
|
|
|
|
|
def _parse_history_date(raw: Any) -> date | None:
|
|
if raw is None:
|
|
return None
|
|
try:
|
|
return date.fromisoformat(str(raw)[:10])
|
|
except ValueError:
|
|
return None
|
|
|
|
|
|
def aggregate_counting_by_date(
|
|
rows: list[dict[str, Any]] | None,
|
|
*,
|
|
kandang_name: str,
|
|
) -> dict[date, int]:
|
|
"""Sum ``total`` for floors whose location starts with ``KN-`` / equals ``KN``."""
|
|
prefix = location_prefix_for_kandang(kandang_name)
|
|
by_date: dict[date, int] = {}
|
|
for row in rows or []:
|
|
day = _parse_history_date(row.get("date"))
|
|
if day is None:
|
|
continue
|
|
location = str(row.get("location") or "").strip()
|
|
if location != prefix and not location.startswith(f"{prefix}-"):
|
|
continue
|
|
by_date[day] = by_date.get(day, 0) + int(row.get("total") or 0)
|
|
return by_date
|
|
|
|
|
|
def aggregate_mortality_by_date(
|
|
rows: list[dict[str, Any]] | None,
|
|
*,
|
|
kandang_name: str,
|
|
) -> dict[date, int]:
|
|
"""Map date → ``total_mortality_count`` for the matching coop."""
|
|
expected = coop_for_kandang(kandang_name).casefold()
|
|
by_date: dict[date, int] = {}
|
|
for row in rows or []:
|
|
day = _parse_history_date(row.get("date"))
|
|
if day is None:
|
|
continue
|
|
coop = str(row.get("coop") or row.get("location") or "").strip()
|
|
if coop.casefold() != expected:
|
|
continue
|
|
by_date[day] = int(row.get("total_mortality_count") or 0)
|
|
return by_date
|
|
|
|
|
|
def cycle_history_window(cycle: Cycle, *, today: date | None = None) -> tuple[date, date]:
|
|
if today is None:
|
|
today = timezone.localdate()
|
|
end = min(today, cycle.end_date)
|
|
return cycle.start_date, end
|
|
|
|
|
|
class ChickenCountingEdgeClient:
|
|
"""Outbound client for the chicken-counting edge box."""
|
|
|
|
DB_HISTORY_PATH = "/api/db/history"
|
|
MORTALITY_HISTORY_PATH = "/api/mortality/history"
|
|
|
|
def __init__(self, base_url: str | None = None, timeout: float | None = None):
|
|
self.base_url = (
|
|
base_url or getattr(settings, "CHICKEN_COUNTING_EDGE_BASE_URL", "") or ""
|
|
).rstrip("/")
|
|
self.timeout = timeout or float(
|
|
getattr(settings, "CHICKEN_COUNTING_EDGE_TIMEOUT_SECONDS", 30) or 30
|
|
)
|
|
|
|
def _get_json(self, path: str) -> Any:
|
|
if not self.base_url:
|
|
raise ChickenCountingEdgeError("CHICKEN_COUNTING_EDGE_BASE_URL is not configured")
|
|
url = urljoin(self.base_url + "/", path.lstrip("/"))
|
|
try:
|
|
with httpx.Client(timeout=self.timeout) as client:
|
|
response = client.get(url)
|
|
except httpx.HTTPError as exc:
|
|
raise ChickenCountingEdgeError(f"chicken counting edge request failed: {exc}") from exc
|
|
if response.status_code >= 400:
|
|
raise ChickenCountingEdgeError(
|
|
f"chicken counting edge HTTP {response.status_code}: {response.text[:300]}",
|
|
status_code=response.status_code,
|
|
)
|
|
return response.json()
|
|
|
|
def fetch_db_history(self) -> list[dict[str, Any]]:
|
|
payload = self._get_json(self.DB_HISTORY_PATH)
|
|
if isinstance(payload, list):
|
|
return payload
|
|
return list(payload.get("results") or payload.get("history") or [])
|
|
|
|
def fetch_mortality_history(self) -> list[dict[str, Any]]:
|
|
payload = self._get_json(self.MORTALITY_HISTORY_PATH)
|
|
if isinstance(payload, list):
|
|
return payload
|
|
return list(payload.get("results") or payload.get("history") or [])
|
|
|
|
|
|
@transaction.atomic
|
|
def sync_chicken_counting_for_cycle(
|
|
cycle: Cycle,
|
|
*,
|
|
db_history: list[dict[str, Any]] | None = None,
|
|
mortality_history: list[dict[str, Any]] | None = None,
|
|
) -> dict[str, Any]:
|
|
"""Upsert ChickenCounting from edge history within the cycle window."""
|
|
counting_enabled = bool(
|
|
getattr(settings, "CHICKEN_COUNTING_EDGE_COUNTING_SYNC_ENABLED", True)
|
|
)
|
|
mortality_enabled = bool(
|
|
getattr(settings, "CHICKEN_COUNTING_EDGE_MORTALITY_SYNC_ENABLED", True)
|
|
)
|
|
if not counting_enabled and not mortality_enabled:
|
|
return {"dates": [], "upserted": 0, "skipped": True}
|
|
|
|
client = ChickenCountingEdgeClient()
|
|
if db_history is None and counting_enabled:
|
|
db_history = client.fetch_db_history()
|
|
if mortality_history is None and mortality_enabled:
|
|
mortality_history = client.fetch_mortality_history()
|
|
|
|
kandang_name = cycle.kandang.kandang_name
|
|
counting_by_date = (
|
|
aggregate_counting_by_date(db_history, kandang_name=kandang_name)
|
|
if counting_enabled
|
|
else {}
|
|
)
|
|
mortality_by_date = (
|
|
aggregate_mortality_by_date(mortality_history, kandang_name=kandang_name)
|
|
if mortality_enabled
|
|
else {}
|
|
)
|
|
|
|
start, end = cycle_history_window(cycle)
|
|
dates = sorted(set(counting_by_date) | set(mortality_by_date))
|
|
upserted: list[date] = []
|
|
|
|
for day in dates:
|
|
if day < start or day > end:
|
|
continue
|
|
defaults: dict[str, int] = {}
|
|
if counting_enabled and day in counting_by_date:
|
|
defaults["total_count"] = counting_by_date[day]
|
|
if mortality_enabled and day in mortality_by_date:
|
|
defaults["mortality_count"] = mortality_by_date[day]
|
|
if not defaults:
|
|
continue
|
|
|
|
obj, created = ChickenCounting.objects.get_or_create(
|
|
cycle=cycle,
|
|
date=day,
|
|
defaults={
|
|
"total_count": defaults.get("total_count", 0),
|
|
"mortality_count": defaults.get("mortality_count", 0),
|
|
},
|
|
)
|
|
if not created:
|
|
update_fields = []
|
|
if "total_count" in defaults and obj.total_count != defaults["total_count"]:
|
|
obj.total_count = defaults["total_count"]
|
|
update_fields.append("total_count")
|
|
if (
|
|
"mortality_count" in defaults
|
|
and obj.mortality_count != defaults["mortality_count"]
|
|
):
|
|
obj.mortality_count = defaults["mortality_count"]
|
|
update_fields.append("mortality_count")
|
|
if update_fields:
|
|
update_fields.append("updated_at")
|
|
obj.save(update_fields=update_fields)
|
|
upserted.append(day)
|
|
|
|
return {"dates": upserted, "upserted": len(upserted), "skipped": False}
|