add chicken counting, mortality api for daily data sync. Also adjust karung sync to sync every 30 sec after publish cutoff time to update real data without clicking sync button manually

This commit is contained in:
Alberto-Audrix committed 2026-09-04 15:32:18 +07:00
1 parent 5ae477a99d
commit 6c21af745d
17 files changed
+791 -119

No files matched your search

@@ -0,0 +1,45 @@
from django.conf import settings
from django.core.management.base import BaseCommand
from apps.farms.models import Cycle
from apps.operations.services.chicken_counting_edge import (
ChickenCountingEdgeError,
sync_chicken_counting_for_cycle,
)
class Command(BaseCommand):
help = "Sync chicken counting + mortality IoT from the edge vision API"
def add_arguments(self, parser):
parser.add_argument("--cycle-id", type=int, default=None)
def handle(self, *args, **options):
counting_on = settings.CHICKEN_COUNTING_EDGE_COUNTING_SYNC_ENABLED
mortality_on = settings.CHICKEN_COUNTING_EDGE_MORTALITY_SYNC_ENABLED
if not counting_on and not mortality_on:
self.stdout.write(
"CHICKEN_COUNTING_EDGE_*_SYNC_ENABLED=false; skipping"
)
return
cycle_id = options.get("cycle_id")
if cycle_id:
cycles = Cycle.objects.filter(pk=cycle_id).select_related("kandang")
else:
cycles = Cycle.objects.filter(status=Cycle.STATUS_ACTIVE).select_related(
"kandang"
)
for cycle in cycles:
try:
summary = sync_chicken_counting_for_cycle(cycle)
except ChickenCountingEdgeError as exc:
self.stderr.write(f"cycle={cycle.pk}: {exc}")
continue
self.stdout.write(
self.style.SUCCESS(
f"cycle={cycle.pk} upserted={summary['upserted']} "
f"dates={','.join(d.isoformat() for d in summary['dates']) or '-'}"
)
)
@@ -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}
@@ -0,0 +1,213 @@
from datetime import date
from unittest.mock import patch
from django.test import TestCase, override_settings
from rest_framework.test import APIClient
from apps.accounts.models import User
from apps.farms.models import Cycle, Kandang, Site
from apps.operations.models import ChickenCounting, ChickenWeight
from apps.operations.services.chicken_counting_edge import (
ChickenCountingEdgeError,
aggregate_counting_by_date,
aggregate_mortality_by_date,
coop_for_kandang,
location_prefix_for_kandang,
sync_chicken_counting_for_cycle,
)
DB_HISTORY_SAMPLE = [
{"date": "2026-09-02", "location": "K1-L1", "total": 14591},
{"date": "2026-09-02", "location": "K1-L2", "total": 14591},
{"date": "2026-09-02", "location": "K2-L1", "total": 100},
{"date": "2026-09-02", "location": "K2-L2", "total": 100},
{"date": "2026-09-03", "location": "K1-L1", "total": 18736},
{"date": "2026-09-03", "location": "K1-L2", "total": 18736},
{"date": "2026-08-20", "location": "K1-L1", "total": 50},
]
MORTALITY_HISTORY_SAMPLE = [
{"date": "2026-09-02", "coop": "K1", "total_mortality_count": 0},
{"date": "2026-09-02", "coop": "K2", "total_mortality_count": 10},
{"date": "2026-09-03", "coop": "K1", "total_mortality_count": 29},
{"date": "2026-09-03", "coop": "K2", "total_mortality_count": 29},
]
class LocationMappingTests(TestCase):
def test_maps_kandang_to_prefix_and_coop(self):
self.assertEqual(location_prefix_for_kandang("Kandang 1"), "K1")
self.assertEqual(coop_for_kandang("Kandang 2"), "K2")
class AggregationTests(TestCase):
def test_sums_matching_floors_per_date(self):
by_date = aggregate_counting_by_date(DB_HISTORY_SAMPLE, kandang_name="Kandang 1")
self.assertEqual(by_date[date(2026, 9, 2)], 29182)
self.assertEqual(by_date[date(2026, 9, 3)], 37472)
self.assertNotIn(date(2026, 9, 2), aggregate_counting_by_date([], kandang_name="Kandang 1"))
def test_mortality_uses_matching_coop(self):
by_date = aggregate_mortality_by_date(
MORTALITY_HISTORY_SAMPLE, kandang_name="Kandang 1"
)
self.assertEqual(by_date[date(2026, 9, 2)], 0)
self.assertEqual(by_date[date(2026, 9, 3)], 29)
@override_settings(
CHICKEN_COUNTING_EDGE_BASE_URL="http://edge.internal",
CHICKEN_COUNTING_EDGE_COUNTING_SYNC_ENABLED=True,
CHICKEN_COUNTING_EDGE_MORTALITY_SYNC_ENABLED=True,
)
class SyncChickenCountingTests(TestCase):
def setUp(self):
user = User.objects.create_user(user_name="sync-user", password="secret")
site = Site.objects.create(site_name="Sukawarna", user=user)
self.kandang = Kandang.objects.create(kandang_name="Kandang 1", site=site)
self.cycle = Cycle.objects.create(
kandang=self.kandang,
total_days=33,
start_date=date(2026, 9, 2),
end_date=date(2026, 10, 4),
doc_in_weight=35,
doc_in_count=50000,
status=Cycle.STATUS_ACTIVE,
)
@patch("apps.operations.services.chicken_counting_edge.timezone.localdate")
@patch("apps.operations.services.chicken_counting_edge.ChickenCountingEdgeClient.fetch_mortality_history")
@patch("apps.operations.services.chicken_counting_edge.ChickenCountingEdgeClient.fetch_db_history")
def test_upserts_within_cycle_window(self, mock_db, mock_mort, mock_today):
mock_today.return_value = date(2026, 9, 4)
mock_db.return_value = DB_HISTORY_SAMPLE
mock_mort.return_value = MORTALITY_HISTORY_SAMPLE
summary = sync_chicken_counting_for_cycle(self.cycle)
self.assertEqual(sorted(summary["dates"]), [date(2026, 9, 2), date(2026, 9, 3)])
row_sep2 = ChickenCounting.objects.get(cycle=self.cycle, date=date(2026, 9, 2))
self.assertEqual(row_sep2.total_count, 29182)
self.assertEqual(row_sep2.mortality_count, 0)
row_sep3 = ChickenCounting.objects.get(cycle=self.cycle, date=date(2026, 9, 3))
self.assertEqual(row_sep3.total_count, 37472)
self.assertEqual(row_sep3.mortality_count, 29)
self.assertFalse(
ChickenCounting.objects.filter(cycle=self.cycle, date=date(2026, 8, 20)).exists()
)
@patch("apps.operations.services.chicken_counting_edge.timezone.localdate")
@patch("apps.operations.services.chicken_counting_edge.ChickenCountingEdgeClient.fetch_mortality_history")
@patch("apps.operations.services.chicken_counting_edge.ChickenCountingEdgeClient.fetch_db_history")
def test_respects_partial_enable_flags(self, mock_db, mock_mort, mock_today):
mock_today.return_value = date(2026, 9, 4)
mock_db.return_value = DB_HISTORY_SAMPLE
mock_mort.return_value = MORTALITY_HISTORY_SAMPLE
ChickenCounting.objects.create(
cycle=self.cycle,
date=date(2026, 9, 2),
total_count=1,
mortality_count=5,
)
with override_settings(
CHICKEN_COUNTING_EDGE_COUNTING_SYNC_ENABLED=False,
CHICKEN_COUNTING_EDGE_MORTALITY_SYNC_ENABLED=True,
):
sync_chicken_counting_for_cycle(self.cycle)
row = ChickenCounting.objects.get(cycle=self.cycle, date=date(2026, 9, 2))
self.assertEqual(row.total_count, 1)
self.assertEqual(row.mortality_count, 0)
@override_settings(
CHICKEN_COUNTING_EDGE_BASE_URL="http://edge.internal",
CHICKEN_COUNTING_EDGE_COUNTING_SYNC_ENABLED=True,
CHICKEN_COUNTING_EDGE_MORTALITY_SYNC_ENABLED=True,
DASHBOARD_PUBLISH_HOUR=17,
TIME_ZONE="Asia/Jakarta",
)
class ChickenCountingSyncApiTests(TestCase):
def setUp(self):
self.user = User.objects.create_user(user_name="tester", password="secret")
self.client = APIClient()
self.client.force_authenticate(self.user)
site = Site.objects.create(site_name="Sukawarna", user=self.user)
self.kandang = Kandang.objects.create(kandang_name="Kandang 1", site=site)
self.cycle = Cycle.objects.create(
kandang=self.kandang,
total_days=33,
start_date=date(2026, 9, 2),
end_date=date(2026, 10, 4),
doc_in_weight=35,
doc_in_count=50000,
status=Cycle.STATUS_ACTIVE,
)
@patch("apps.operations.services.chicken_counting_edge.timezone.localdate")
@patch("apps.operations.services.chicken_counting_edge.ChickenCountingEdgeClient.fetch_mortality_history")
@patch("apps.operations.services.chicken_counting_edge.ChickenCountingEdgeClient.fetch_db_history")
def test_sync_action_upserts(self, mock_db, mock_mort, mock_today):
mock_today.return_value = date(2026, 9, 4)
mock_db.return_value = DB_HISTORY_SAMPLE
mock_mort.return_value = MORTALITY_HISTORY_SAMPLE
response = self.client.post(
"/api/v1/chicken-countings/sync/",
{"cycle_id": self.cycle.pk},
format="json",
)
self.assertEqual(response.status_code, 200, response.data)
self.assertEqual(ChickenCounting.objects.filter(cycle=self.cycle).count(), 2)
@patch(
"apps.operations.views.sync_chicken_counting_for_cycle",
side_effect=ChickenCountingEdgeError("edge down"),
)
def test_sync_action_returns_502_on_edge_error(self, _mock_sync):
response = self.client.post(
"/api/v1/chicken-countings/sync/",
{"cycle_id": self.cycle.pk},
format="json",
)
self.assertEqual(response.status_code, 502)
self.assertIn("edge down", response.data["detail"])
def test_counting_and_weight_visible_before_publish_time(self):
ChickenCounting.objects.create(
cycle=self.cycle,
date=date(2026, 9, 4),
total_count=100,
mortality_count=2,
)
ChickenWeight.objects.create(
cycle=self.cycle,
date=date(2026, 9, 4),
age=2,
doc_weight=35,
average_weight=40,
chicken_count=10,
uniformity=90,
average_daily_gain=1,
)
with patch("apps.operations.services.visibility.timezone.localtime") as mock_localtime:
from datetime import datetime
from zoneinfo import ZoneInfo
wib = ZoneInfo("Asia/Jakarta")
mock_localtime.return_value = datetime(2026, 9, 4, 10, 0, tzinfo=wib)
counting = self.client.get(
"/api/v1/chicken-countings/", {"cycle_id": self.cycle.pk}
)
weight = self.client.get(
"/api/v1/chicken-weights/", {"cycle_id": self.cycle.pk}
)
self.assertEqual(counting.status_code, 200)
c_results = counting.data["results"] if isinstance(counting.data, dict) else counting.data
self.assertEqual(len(c_results), 1)
self.assertEqual(c_results[0]["date"], "2026-09-04")
self.assertEqual(weight.status_code, 200)
w_results = weight.data["results"] if isinstance(weight.data, dict) else weight.data
self.assertEqual(len(w_results), 1)
self.assertEqual(w_results[0]["date"], "2026-09-04")
+41
View File
@@ -28,6 +28,10 @@ from apps.operations.serializers import (
ManualInputSerializer,
)
from apps.operations.services.karung_web import KarungWebError, sync_karung_for_cycle
from apps.operations.services.chicken_counting_edge import (
ChickenCountingEdgeError,
sync_chicken_counting_for_cycle,
)
from apps.operations.services.visibility import dashboard_publish_time, visible_through_date
READ_ACTIONS = frozenset({"list", "retrieve", "latest_average", "latest", "dates"})
@@ -70,14 +74,51 @@ class ManualInputViewSet(CycleScopedViewSet):
class ChickenCountingViewSet(CycleScopedViewSet):
"""Chicken counting + mortality IoT — not subject to dashboard publish cutoff."""
queryset = ChickenCounting.objects.select_related("cycle").all()
serializer_class = ChickenCountingSerializer
def applies_visibility_filter(self) -> bool:
return False
@action(detail=False, methods=["post"], url_path="sync")
def sync_from_edge(self, request):
cycle_id = request.data.get("cycle_id")
if not cycle_id:
return Response(
{"detail": "cycle_id is required."},
status=status.HTTP_400_BAD_REQUEST,
)
try:
cycle = Cycle.objects.select_related("kandang").get(pk=cycle_id)
except Cycle.DoesNotExist:
return Response({"detail": "Cycle not found."}, status=status.HTTP_404_NOT_FOUND)
try:
summary = sync_chicken_counting_for_cycle(cycle)
except (ChickenCountingEdgeError, ValidationError) as exc:
code = getattr(exc, "status_code", None)
return Response(
{"detail": str(exc)},
status=code if code and code >= 400 else status.HTTP_502_BAD_GATEWAY,
)
return Response(
{
"upserted": summary["upserted"],
"dates": [d.isoformat() for d in summary["dates"]],
}
)
class ChickenWeightViewSet(CycleScopedViewSet):
"""Chicken weight — not subject to dashboard publish cutoff."""
queryset = ChickenWeight.objects.select_related("cycle").all()
serializer_class = ChickenWeightSerializer
def applies_visibility_filter(self) -> bool:
return False
@action(detail=False, methods=["get"], url_path="latest-avg")
def latest_average(self, request):
cycle_id = request.query_params.get("cycle_id")
+14
View File
@@ -195,3 +195,17 @@ IOT_SYNC_LOOKBACK_MINUTES = int(env("IOT_SYNC_LOOKBACK_MINUTES", "20") or "20")
IOT_SYNC_MAX_PAGES = int(env("IOT_SYNC_MAX_PAGES", "10") or "10")
# JSON map: external IoT flock_id -> Django flock.pk, e.g. {"686df69f407b21002da8750c": 1}
IOT_FLOCK_ID_MAP = env_json_dict("IOT_FLOCK_ID_MAP")
# Chicken counting edge vision API (headcount + mortality) — not gated by publish cutoff.
CHICKEN_COUNTING_EDGE_BASE_URL = (
env("CHICKEN_COUNTING_EDGE_BASE_URL", "") or ""
).rstrip("/")
CHICKEN_COUNTING_EDGE_TIMEOUT_SECONDS = float(
env("CHICKEN_COUNTING_EDGE_TIMEOUT_SECONDS", "30") or "30"
)
CHICKEN_COUNTING_EDGE_COUNTING_SYNC_ENABLED = env_bool(
"CHICKEN_COUNTING_EDGE_COUNTING_SYNC_ENABLED", True
)
CHICKEN_COUNTING_EDGE_MORTALITY_SYNC_ENABLED = env_bool(
"CHICKEN_COUNTING_EDGE_MORTALITY_SYNC_ENABLED", True
)