update iot panel will now listen to push from panel
This commit is contained in:
1 parent
20d04acb3e
commit
520c4bcaad
10 files changed
+158
-754
No files matched your search
@@ -1,14 +1,10 @@
|
||||
"""Sync IoT panel readings from dashboard.cpsp.id/api/iot/flocks/ (10-minute snapshots)."""
|
||||
"""Handle IoT panel push readings (webhook from panel IoT devices)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from datetime import date, datetime, timedelta, timezone as dt_timezone
|
||||
from typing import Any
|
||||
from datetime import date, datetime, timezone as dt_timezone
|
||||
from zoneinfo import ZoneInfo
|
||||
|
||||
import httpx
|
||||
from django.conf import settings
|
||||
from django.utils import timezone
|
||||
|
||||
from apps.farms.models import Cycle, Flock
|
||||
@@ -28,65 +24,15 @@ CHILL_FACTOR_POINTS = (
|
||||
)
|
||||
|
||||
|
||||
class IotApiError(Exception):
|
||||
class IotPushError(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 _to_float(value) -> float:
|
||||
raw = float(value)
|
||||
return raw / 10.0 if raw > 100 else raw
|
||||
|
||||
|
||||
def _chicken_chill_factor(age: float) -> float:
|
||||
@@ -118,174 +64,52 @@ def _experience_temperature(
|
||||
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
|
||||
def handle_iot_push(flock: Flock, payload: dict) -> IotPanel:
|
||||
now = timezone.now()
|
||||
snapshot_date = now.astimezone(WIB).date()
|
||||
|
||||
avg_temp = _to_float(payload.get("suhu", 0))
|
||||
humidity = _to_float(payload.get("kelembapan", 0))
|
||||
wind_speed = _to_float(payload.get("wind", 0))
|
||||
water_total = float(payload.get("flow", 0) or 0)
|
||||
|
||||
class IotApiClient:
|
||||
def __init__(self, base_url: str, timeout: float | None = None):
|
||||
self.base_url = (base_url or "").strip()
|
||||
if not self.base_url:
|
||||
raise IotApiError("IoT API URL is empty")
|
||||
self.timeout = timeout or settings.IOT_API_SYNC_TIMEOUT_SECONDS
|
||||
|
||||
def fetch_page(self, page: int) -> dict[str, Any]:
|
||||
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
|
||||
|
||||
parts = urlsplit(self.base_url)
|
||||
query = dict(parse_qsl(parts.query, keep_blank_values=True))
|
||||
query["page"] = str(page)
|
||||
url = urlunsplit(
|
||||
(parts.scheme, parts.netloc, parts.path, urlencode(query), parts.fragment)
|
||||
cycle = (
|
||||
Cycle.objects.filter(
|
||||
kandang_id=flock.kandang_id,
|
||||
start_date__lte=snapshot_date,
|
||||
end_date__gte=snapshot_date,
|
||||
)
|
||||
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()
|
||||
.order_by("-start_date")
|
||||
.first()
|
||||
)
|
||||
age = cycle_age_for_date(cycle.start_date, snapshot_date) if cycle else 1
|
||||
|
||||
experience_temp = round(
|
||||
_experience_temperature(
|
||||
average_temperature=avg_temp,
|
||||
wind_speed=wind_speed,
|
||||
humidity_pct=humidity,
|
||||
age_days=age,
|
||||
),
|
||||
4,
|
||||
)
|
||||
|
||||
def _sync_flock_pages(
|
||||
flock: Flock,
|
||||
*,
|
||||
lookback_minutes: int,
|
||||
max_pages: int,
|
||||
cutoff,
|
||||
) -> SyncSummary:
|
||||
client = IotApiClient(flock.iot_api_url)
|
||||
upserted = 0
|
||||
skipped = 0
|
||||
pages_read = 0
|
||||
|
||||
for page in range(1, max_pages + 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:
|
||||
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)
|
||||
|
||||
|
||||
def sync_iot_panels(
|
||||
*,
|
||||
flock_id: int | None = None,
|
||||
lookback_minutes: int | None = None,
|
||||
max_pages: int | None = None,
|
||||
) -> SyncSummary:
|
||||
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)
|
||||
|
||||
flocks = Flock.objects.select_related("kandang").exclude(iot_api_url="")
|
||||
if flock_id is not None:
|
||||
flocks = flocks.filter(pk=flock_id)
|
||||
|
||||
flock_list = list(flocks)
|
||||
if not flock_list:
|
||||
return SyncSummary(upserted=0, skipped=0, pages=0)
|
||||
|
||||
upserted = 0
|
||||
skipped = 0
|
||||
pages_read = 0
|
||||
for flock in flock_list:
|
||||
if not (flock.iot_api_url or "").strip():
|
||||
continue
|
||||
summary = _sync_flock_pages(
|
||||
flock,
|
||||
lookback_minutes=lookback,
|
||||
max_pages=page_limit,
|
||||
cutoff=cutoff,
|
||||
)
|
||||
upserted += summary.upserted
|
||||
skipped += summary.skipped
|
||||
pages_read += summary.pages
|
||||
|
||||
return SyncSummary(upserted=upserted, skipped=skipped, pages=pages_read)
|
||||
panel, _created = IotPanel.objects.update_or_create(
|
||||
flock=flock,
|
||||
timestamp=now,
|
||||
defaults={
|
||||
"date": snapshot_date,
|
||||
"wind_speed": round(wind_speed, 4),
|
||||
"humidity": round(humidity, 4),
|
||||
"water_total": water_total,
|
||||
"average_temperature": round(avg_temp, 4),
|
||||
"experience_temperature": experience_temp,
|
||||
},
|
||||
)
|
||||
return panel
|
||||
Reference in new issue
Block a user