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

335 lines
11 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.utils import timezone
from apps.farms.models import Cycle
from apps.operations.models import FeedSacks
class KarungWebError(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:
"""Extract the numeric kandang index from labels like ``Kandang 1``."""
match = re.search(r"(\d+)", kandang_name or "")
if not match:
raise ValidationError(
f"Cannot derive IoT kandang index from kandang name: {kandang_name!r}"
)
return match.group(1)
def source_for_kandang(kandang_name: str, *, kind: str) -> str:
"""Map a kandang label to a karung-web-admin source such as ``K1 - In``."""
return f"K{kandang_index(kandang_name)} - {kind}"
def masuk_source_for_kandang(kandang_name: str) -> str:
"""Map a kandang label like ``Kandang 1`` to the karung-web-admin source ``K1 - In``."""
return source_for_kandang(kandang_name, kind="In")
def _karung_for_source(rows: list[dict[str, Any]] | None, expected_source: str) -> int:
normalized = _normalize_source(expected_source)
for row in rows or []:
if _normalize_source(row.get("source")) == normalized:
return int(row.get("karung") or row.get("counter_value") or 0)
return 0
def _normalize_source(value: str | None) -> str:
return (value or "").strip().casefold()
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 daily_totals_from_combined_history(
combined: dict[str, Any],
*,
kandang_name: str,
) -> dict[date, dict[str, int]]:
"""Group /api/combined history into per-date in/out/use for one kandang."""
in_source = _normalize_source(source_for_kandang(kandang_name, kind="In"))
out_source = _normalize_source(source_for_kandang(kandang_name, kind="Out"))
use_source = _normalize_source(source_for_kandang(kandang_name, kind="Use"))
by_date: dict[date, dict[str, int]] = {}
def ensure(day: date) -> dict[str, int]:
return by_date.setdefault(
day, {"in_today": 0, "out_today": 0, "feed_use_today": 0}
)
masuk_history = (combined.get("masuk") or {}).get("history") or []
for row in masuk_history:
day = _parse_history_date(row.get("date"))
if day is None:
continue
source = _normalize_source(row.get("source"))
value = int(row.get("counter_value") or row.get("karung") or 0)
if source == in_source:
ensure(day)["in_today"] = value
elif source == out_source:
ensure(day)["out_today"] = value
tuang_history = (combined.get("tuang") or {}).get("history") or []
for row in tuang_history:
day = _parse_history_date(row.get("date"))
if day is None:
continue
source = _normalize_source(row.get("source"))
value = int(row.get("counter_value") or row.get("karung") or 0)
if source == use_source:
ensure(day)["feed_use_today"] = value
return by_date
def cycle_history_window(cycle: Cycle, *, today: date | None = None) -> tuple[date, date]:
"""Inclusive date window for history upserts (calendar today while open)."""
if today is None:
today = timezone.localdate()
end = cycle.operational_end_date(today=today)
return cycle.start_date, end
def backfill_karung_from_combined(
cycle: Cycle,
combined: dict[str, Any] | None = None,
) -> list[date]:
"""Upsert FeedSacks for cycle dates present in combined history. Returns upserted dates."""
if combined is None:
combined = KarungWebClient().fetch_combined()
kandang_name = cycle.kandang.kandang_name
by_date = daily_totals_from_combined_history(combined, kandang_name=kandang_name)
start, end = cycle_history_window(cycle)
upserted: list[date] = []
for day in sorted(by_date):
if day < start or day > end:
continue
totals = by_date[day]
FeedSacks.objects.update_or_create(
cycle=cycle,
date=day,
defaults={
"in_today": totals["in_today"],
"out_today": totals["out_today"],
"feed_use_today": totals["feed_use_today"],
},
)
upserted.append(day)
return upserted
def recompute_feed_sack_totals(cycle: Cycle) -> None:
"""Recompute running in/out/use totals for all FeedSacks of a cycle in date order."""
running_in = 0
running_out = 0
running_use = 0
rows = list(FeedSacks.objects.filter(cycle=cycle).order_by("date", "pk"))
for row in rows:
running_in += row.in_today or 0
running_out += row.out_today or 0
running_use += row.feed_use_today or 0
changed = (
row.in_total != running_in
or row.out_total != running_out
or row.feed_use_total != running_use
)
if changed:
row.in_total = running_in
row.out_total = running_out
row.feed_use_total = running_use
row.save(update_fields=["in_total", "out_total", "feed_use_total", "updated_at"])
class KarungWebClient:
"""Outbound client for karung-web-admin — base URL from settings only."""
SECTIONS_PATH = "/api/v1/cycles/active/sections"
COMBINED_PATH = "/api/combined"
def __init__(self, base_url: str | None = None, timeout: float | None = None):
self.base_url = (base_url or settings.KARUNG_WEB_ADMIN_BASE_URL).rstrip("/")
self.timeout = timeout or settings.KARUNG_WEB_ADMIN_TIMEOUT_SECONDS
def _get_json(self, path: str) -> dict[str, Any]:
if not self.base_url:
raise KarungWebError("KARUNG_WEB_ADMIN_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 KarungWebError(f"karung-web-admin request failed: {exc}") from exc
if response.status_code >= 400:
raise KarungWebError(
f"karung-web-admin HTTP {response.status_code}: {response.text[:300]}",
status_code=response.status_code,
)
payload = response.json()
if not payload.get("success", True) and payload.get("error"):
error = payload.get("error") or {}
raise KarungWebError(error.get("message") or "karung-web-admin returned failure")
return payload.get("data") or payload
def fetch_active_sections(self) -> dict[str, Any]:
return self._get_json(self.SECTIONS_PATH)
def fetch_combined(self) -> dict[str, Any]:
return self._get_json(self.COMBINED_PATH)
def fetch_iot_masuk_for_date(target_date: date, *, kandang_name: str) -> int | None:
"""Return IoT karung-masuk for ``target_date`` and kandang from karung-web-admin /api/combined."""
expected_source = _normalize_source(masuk_source_for_kandang(kandang_name))
target = target_date.isoformat()
client = KarungWebClient()
data = client.fetch_combined()
masuk = data.get("masuk") or {}
for row in masuk.get("history") or []:
if row.get("date") == target and _normalize_source(row.get("source")) == expected_source:
return int(row.get("counter_value") or 0)
if target_date == date.today():
for row in masuk.get("today") or []:
if _normalize_source(row.get("source")) == expected_source:
return int(row.get("karung") or 0)
return None
def initial_balance_accuracy(manual: int, iot_in: int | None) -> float | None:
if manual > 0 and iot_in is not None and iot_in > 0:
return round(min(manual, iot_in) / max(manual, iot_in) * 1000) / 10
if manual == 0 and iot_in == 0:
return 100.0
return None
def pull_initial_balance_iot(cycle: Cycle, balance_date: date) -> tuple[int | None, str | None]:
sync_error = None
iot_in = None
try:
iot_in = fetch_iot_masuk_for_date(
balance_date,
kandang_name=cycle.kandang.kandang_name,
)
except KarungWebError as exc:
sync_error = str(exc)
return iot_in, sync_error
def request_karung(
cycle: Cycle,
*,
section: str = "today",
) -> tuple[FeedSacks, dict[str, Any]]:
if section not in {"today", "yesterday_history"}:
raise ValidationError("section must be today or yesterday_history")
client = KarungWebClient()
data = client.fetch_active_sections()
sections = data.get("sections") or {}
block = sections.get(section) or {}
totals = block.get("totals") or {}
remote_cycle = data.get("cycle") or {}
kandang_name = cycle.kandang.kandang_name
k_index = kandang_index(kandang_name)
masuk_rows = block.get("masuk") or []
tuang_rows = block.get("tuang") or []
in_today = _karung_for_source(masuk_rows, f"K{k_index} - In")
feed_use_today = _karung_for_source(tuang_rows, f"K{k_index} - Use")
out_today = _karung_for_source(masuk_rows, f"K{k_index} - Out")
saldo_awal = int(remote_cycle.get("saldo_awal") or 0)
snapshot_raw = block.get("date")
snapshot_date = date.fromisoformat(str(snapshot_raw)) if snapshot_raw else timezone.localdate()
prior = (
FeedSacks.objects.filter(cycle=cycle)
.exclude(date=snapshot_date)
.order_by("-date", "-pk")
.first()
)
if prior:
in_total = prior.in_total + in_today
out_total = prior.out_total + out_today
feed_use_total = prior.feed_use_total + feed_use_today
else:
in_total = in_today
out_total = out_today
feed_use_total = feed_use_today
feed_sack, _ = FeedSacks.objects.update_or_create(
cycle=cycle,
date=snapshot_date,
defaults={
"in_today": in_today,
"out_today": out_today,
"in_total": in_total,
"out_total": out_total,
"feed_use_today": feed_use_today,
"feed_use_total": feed_use_total,
},
)
summary = {
"section": section,
"snapshot_date": snapshot_date.isoformat(),
"totals": totals,
"kandang": kandang_name,
"kandang_sources": {
"in": source_for_kandang(kandang_name, kind="In"),
"use": source_for_kandang(kandang_name, kind="Use"),
"out": source_for_kandang(kandang_name, kind="Out"),
},
"kandang_totals": {
"in_today": in_today,
"feed_use_today": feed_use_today,
"out_today": out_today,
},
"saldo_awal": saldo_awal,
"remote_cycle": remote_cycle,
}
return feed_sack, summary
def sync_karung_for_cycle(
cycle: Cycle,
*,
section: str = "today",
) -> tuple[FeedSacks, dict[str, Any]]:
"""Sync today from sections, backfill cycle history from /api/combined, recompute totals."""
feed_sack, summary = request_karung(cycle, section=section)
history_dates = backfill_karung_from_combined(cycle)
recompute_feed_sack_totals(cycle)
feed_sack.refresh_from_db()
summary = {
**summary,
"history_dates": [d.isoformat() for d in history_dates],
}
return feed_sack, summary