Files
dashboard-cpsp/backend/apps/operations/services/chicken_counting_edge.py
T

333 lines
12 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, ChickenWeight
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 aggregate_weight_by_date(
rows: list[dict[str, Any]] | None,
*,
kandang_name: str,
) -> dict[date, dict[str, float | int]]:
"""Map date → fused weight metrics for the matching coop.
Prefer ``camera_id=FUSED``. When both floors (e.g. K1-L1 / K1-L2) carry the
same fused prediction, average weight/uniformity and take max sample count
so duplicates are not double-counted.
"""
expected = coop_for_kandang(kandang_name).casefold()
floors_by_date: dict[date, list[dict[str, Any]]] = {}
for row in rows or []:
if str(row.get("camera_id") or "").strip().upper() != "FUSED":
continue
day = _parse_history_date(row.get("date"))
if day is None:
continue
coop = str(row.get("coop") or "").strip()
if coop.casefold() != expected:
continue
floors_by_date.setdefault(day, []).append(row)
by_date: dict[date, dict[str, float | int]] = {}
for day, floors in floors_by_date.items():
weights = [float(f.get("predicted_weight_g") or 0) for f in floors]
unifs = [float(f.get("uniformity") or 0) for f in floors]
samples = [int(f.get("n_valid_samples") or 0) for f in floors]
by_date[day] = {
"average_weight": sum(weights) / len(weights),
"uniformity": sum(unifs) / len(unifs),
"chicken_count": max(samples) if samples else 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 = cycle.operational_end_date(today=today)
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"
WEIGHT_PREDICTIONS_PATH = "/api/weight/predictions"
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 [])
def fetch_weight_predictions(self) -> list[dict[str, Any]]:
payload = self._get_json(self.WEIGHT_PREDICTIONS_PATH)
if isinstance(payload, list):
return payload
return list(payload.get("predictions") or payload.get("results") 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}
@transaction.atomic
def sync_chicken_weight_for_cycle(
cycle: Cycle,
*,
predictions: list[dict[str, Any]] | None = None,
) -> dict[str, Any]:
"""Upsert ChickenWeight from edge weight predictions within the cycle window.
Day 0 (cycle start) always uses DOC weight and 0% uniformity — edge values
for that date are ignored.
"""
if not bool(getattr(settings, "CHICKEN_COUNTING_EDGE_WEIGHT_SYNC_ENABLED", True)):
return {"dates": [], "upserted": 0, "skipped": True}
client = ChickenCountingEdgeClient()
if predictions is None:
predictions = client.fetch_weight_predictions()
kandang_name = cycle.kandang.kandang_name
weight_by_date = aggregate_weight_by_date(predictions, kandang_name=kandang_name)
start, end = cycle_history_window(cycle)
if start > end:
return {"dates": [], "upserted": 0, "skipped": False}
doc_weight = int(cycle.doc_in_weight)
dates = sorted({d for d in weight_by_date if start <= d <= end} | {start})
prev_weight = float(doc_weight)
# Seed previous-day weight from any row before the first upsert date.
existing_before = (
ChickenWeight.objects.filter(cycle=cycle, date__lt=dates[0])
.order_by("-date")
.first()
)
if existing_before is not None:
prev_weight = float(existing_before.average_weight)
upserted: list[date] = []
for day in dates:
age = (day - start).days
if age == 0:
average_weight = float(doc_weight)
uniformity = 0.0
chicken_count = int(cycle.doc_in_count)
adg = 0.0
else:
metrics = weight_by_date.get(day)
if metrics is None:
continue
average_weight = float(metrics["average_weight"])
uniformity = float(metrics["uniformity"])
chicken_count = int(metrics["chicken_count"])
adg = round(average_weight - prev_weight, 2)
defaults = {
"age": age,
"doc_weight": doc_weight,
"average_weight": average_weight,
"chicken_count": chicken_count,
"uniformity": uniformity,
"average_daily_gain": adg,
}
ChickenWeight.objects.update_or_create(
cycle=cycle,
date=day,
defaults=defaults,
)
upserted.append(day)
prev_weight = average_weight
return {"dates": upserted, "upserted": len(upserted), "skipped": False}