"""Sync IoT panel readings from dashboard.cpsp.id/api/iot/flocks/ (10-minute snapshots).""" from __future__ import annotations from dataclasses import dataclass from datetime import date, datetime, timedelta, timezone as dt_timezone from typing import Any from zoneinfo import ZoneInfo import httpx from django.conf import settings from django.utils import timezone from apps.farms.models import Cycle, Flock from apps.operations.models import IotPanel WIB = ZoneInfo("Asia/Jakarta") CHILL_FACTOR_POINTS = ( (0, 8.0), (1, 8.0), (7, 7.0), (14, 6.0), (21, 4.5), (28, 3.5), (35, 3.5), (42, 3.0), ) class IotApiError(Exception): def __init__(self, message: str, status_code: int | None = None): super().__init__(message) self.status_code = status_code def _sensor_celsius(sensor: dict | float | int | None) -> float | None: if sensor is None: return None if isinstance(sensor, (int, float)): return float(sensor) raw = sensor.get("value") if raw is None: return None cal = float(sensor.get("calibration") or 0) return (float(raw) + cal) / 10.0 def _sensor_humidity_pct(sensor: dict | float | int | None, fallback: float | None) -> float: if isinstance(sensor, dict): raw = sensor.get("value") if raw is not None: cal = float(sensor.get("calibration") or 0) return (float(raw) + cal) / 10.0 if isinstance(sensor, (int, float)): return float(sensor) if fallback is None: return 70.0 value = float(fallback) return value / 10.0 if value > 100 else value def _wind_ms(payload_wind: float | int | None, sensors_wind: dict | float | int | None) -> float: if isinstance(sensors_wind, dict) and sensors_wind.get("value") is not None: return float(sensors_wind["value"]) / 10.0 if isinstance(sensors_wind, (int, float)): return float(sensors_wind) if payload_wind is None: return 0.0 return float(payload_wind) / 10.0 def _average_inside_temp_c(sensors: dict | None, actual_temperature: float | int | None) -> float: if sensors: readings = [ v for v in ( _sensor_celsius(sensors.get("temperature1")), _sensor_celsius(sensors.get("temperature2")), _sensor_celsius(sensors.get("temperature3")), ) if v is not None ] if readings: return sum(readings) / len(readings) if actual_temperature is None: return 0.0 value = float(actual_temperature) return value / 10.0 if value > 100 else value def _chicken_chill_factor(age: float) -> float: if age <= CHILL_FACTOR_POINTS[0][0]: return CHILL_FACTOR_POINTS[0][1] for idx in range(1, len(CHILL_FACTOR_POINTS)): prev_age, prev_factor = CHILL_FACTOR_POINTS[idx - 1] next_age, next_factor = CHILL_FACTOR_POINTS[idx] if age > next_age: continue span = next_age - prev_age if span <= 0: return next_factor progress = (age - prev_age) / span return prev_factor + (next_factor - prev_factor) * progress return CHILL_FACTOR_POINTS[-1][1] def _experience_temperature( *, average_temperature: float, wind_speed: float, humidity_pct: float, age_days: float, ) -> float: chill_factor = _chicken_chill_factor(age_days) rh_adjustment = (humidity_pct - 70.0) / 5.0 chill_effect = wind_speed * chill_factor - rh_adjustment return average_temperature - chill_effect def panel_fields_from_payload(data: dict, age_days: int) -> dict[str, float]: sensors = data.get("sensors") or {} humidity = _sensor_humidity_pct(sensors.get("humidity"), data.get("humidity")) wind_speed = _wind_ms(data.get("wind"), sensors.get("wind")) avg_temp = _average_inside_temp_c(sensors, data.get("actualTemperature")) water = data.get("water") if water is None and isinstance(sensors.get("water"), dict): water = sensors["water"].get("value") return { "wind_speed": round(wind_speed, 4), "humidity": round(humidity, 4), "water_total": float(water or 0), "average_temperature": round(avg_temp, 4), "experience_temperature": round( _experience_temperature( average_temperature=avg_temp, wind_speed=wind_speed, humidity_pct=humidity, age_days=age_days, ), 4, ), } def parse_fetched_at(value: str) -> datetime: parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) if timezone.is_naive(parsed): return timezone.make_aware(parsed, dt_timezone.utc) return parsed def cycle_age_for_date(cycle_start: date, snapshot_date: date) -> int: if snapshot_date < cycle_start: return 1 return (snapshot_date - cycle_start).days + 1 @dataclass class SyncSummary: upserted: int skipped: int pages: int class IotApiClient: def __init__(self, base_url: str | None = None, timeout: float | None = None): self.base_url = (base_url or settings.IOT_API_BASE_URL).rstrip("/") + "/" self.timeout = timeout or settings.IOT_API_SYNC_TIMEOUT_SECONDS def fetch_page(self, page: int) -> dict[str, Any]: url = f"{self.base_url}?page={page}" try: with httpx.Client(timeout=self.timeout) as client: response = client.get(url) except httpx.HTTPError as exc: raise IotApiError(f"IoT API request failed: {exc}") from exc if response.status_code >= 400: raise IotApiError( f"IoT API HTTP {response.status_code}: {response.text[:300]}", status_code=response.status_code, ) return response.json() def sync_iot_panels( *, lookback_minutes: int | None = None, max_pages: int | None = None, ) -> SyncSummary: flock_map: dict[str, int] = getattr(settings, "IOT_FLOCK_ID_MAP", {}) or {} if not flock_map: return SyncSummary(upserted=0, skipped=0, pages=0) lookback = lookback_minutes if lookback_minutes is not None else settings.IOT_SYNC_LOOKBACK_MINUTES page_limit = max_pages if max_pages is not None else settings.IOT_SYNC_MAX_PAGES cutoff = timezone.now() - timedelta(minutes=lookback) django_flocks: dict[str, Flock] = {} for external_id, django_id in flock_map.items(): try: django_flocks[str(external_id)] = Flock.objects.select_related("kandang").get(pk=int(django_id)) except Flock.DoesNotExist: continue client = IotApiClient() upserted = 0 skipped = 0 pages_read = 0 for page in range(1, page_limit + 1): payload = client.fetch_page(page) pages_read += 1 results = payload.get("results") or [] if not results: break page_has_recent = False for record in results: external_flock_id = str(record.get("flock_id") or "") flock = django_flocks.get(external_flock_id) if flock is None: skipped += 1 continue fetched_raw = record.get("fetched_at") payload_data = (record.get("payload") or {}).get("data") if not fetched_raw or not payload_data: skipped += 1 continue fetched_at = parse_fetched_at(fetched_raw) if fetched_at >= cutoff: page_has_recent = True elif page > 1: continue snapshot_date = fetched_at.astimezone(WIB).date() cycle = ( Cycle.objects.filter( kandang_id=flock.kandang_id, start_date__lte=snapshot_date, end_date__gte=snapshot_date, ) .order_by("-start_date") .first() ) age = cycle_age_for_date(cycle.start_date, snapshot_date) if cycle else 1 fields = panel_fields_from_payload(payload_data, age) IotPanel.objects.update_or_create( flock=flock, timestamp=fetched_at, defaults={ "date": snapshot_date, **fields, }, ) upserted += 1 if not page_has_recent and page > 1: break return SyncSummary(upserted=upserted, skipped=skipped, pages=pages_read)