first commit
This commit is contained in:
commit
b54624be96
226 files changed
+108840
No files matched your search
@@ -0,0 +1,214 @@
|
||||
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}
|
||||
Reference in new issue
Block a user