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

982 lines
35 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Unified AI Insight generate pipeline: grade → RAG → narrate → cache.
Called from `AIInsightViewSet` (`POST generate/`, `GET cached/`).
Numbers/status come from `cp707_knowledge`; Ollama only narrates.
Architecture map: docs/ai-insight/README.md
"""
from __future__ import annotations
import json
import logging
import re
from typing import Any
import httpx
from django.conf import settings
from django.utils import timezone
from apps.farms.models import Cycle, Kandang
from apps.operations.models import AIInsight
from apps.operations.services import cp707_knowledge as cp707
from apps.operations.services import root_cause
logger = logging.getLogger(__name__)
_NUMERIC_PIPE_ROW = re.compile(r"^[\d.,]+(\s*\|\s*[\d.,]*)+$")
# Stored in AIInsight.alert — condition of graded data, not narrative text.
ALERT_HEALTHY = "healthy"
ALERT_WARNING = "warning"
ALERT_CRITICAL = "critical"
ALERT_UNKNOWN = "unknown"
_ALERT_RANK = {
ALERT_UNKNOWN: 0,
ALERT_HEALTHY: 1,
ALERT_WARNING: 2,
ALERT_CRITICAL: 3,
}
_GRADER_STATUS_TO_ALERT = {
"ok": ALERT_HEALTHY,
"healthy": ALERT_HEALTHY,
"warning": ALERT_WARNING,
"critical": ALERT_CRITICAL,
"unknown": ALERT_UNKNOWN,
}
def alert_from_graded(graded: dict[str, Any] | None) -> str:
"""Worst condition among graded analysis blocks (critical > warning > healthy)."""
analyses = (graded or {}).get("analyses") or {}
if not isinstance(analyses, dict) or not analyses:
return ALERT_UNKNOWN
worst = ALERT_UNKNOWN
saw_status = False
for block in analyses.values():
if not isinstance(block, dict):
continue
raw = block.get("status")
if raw is None:
continue
mapped = _GRADER_STATUS_TO_ALERT.get(str(raw).strip().lower())
if mapped is None:
continue
saw_status = True
if _ALERT_RANK[mapped] > _ALERT_RANK[worst]:
worst = mapped
return worst if saw_status else ALERT_UNKNOWN
ANTI_HALLUCINATION_RULES = """
ATURAN NAMA KANDANG & ID:
- Pakai nama kandang PERSIS seperti di Data Halaman / prompt (mis. "Kandang 2").
Jangan menambah kata "Kandang" lagi (hindari "Kandang Kandang 2").
- JANGAN menyebut id kandang, cycle id, atau "(id=…)".
ATURAN MORTALITAS (WAJIB DIPATUHI — TIDAK BOLEH DILANGGAR):
- Bedakan DUA metrik berbeda:
1) mortalitas_hari_ini_ekor = jumlah EKOR mati hari ini (bukan persen).
Jika nilainya 0, tulis mortalitas hari ini 0 ekor — JANGAN sebut "5%" atau "TINGGI".
2) mortalitas kumulatif % = hanya dari field *_persen / persen_hidup / [GRADED FACTS].
- Ambang CP 707 untuk mortalitas KUMULATIF %: <5% normal; 5–7% TINGGI; >7% SANGAT TINGGI.
Ambang ini BUKAN untuk mortalitas hari ini.
- WAJIB: untuk mortalitas kumulatif, sebutkan severity_label dari [GRADED FACTS]
PERSIS (contoh: "SANGAT TINGGI"). Jangan mengganti dengan sinonim lain.
- DILARANG memakai kata di dilarang_kata pada blok mortality (mis. "rendah", "normal",
"baik") jika severity_label = TINGGI atau SANGAT TINGGI.
- DILARANG KERAS menulis kalimat kontradiktif seperti "mortalitas rendah … 9%" atau
"mortalitas normal melebihi batas kritis". Satu arah saja: ikuti status + severity_label.
ATURAN ANGKA & ARAH (WAJIB DIPATUHI):
- Angka di blok STANDAR CP 707 adalah STANDAR saja — BUKAN nilai aktual kandang.
Jangan menulis standar (mis. bobot 2500 g atau mortalitas 5.75%) seolah data aktual.
- Nilai aktual HANYA dari [Data Halaman (JSON)] atau [GRADED FACTS] yang berstatus bukan unknown.
- Jika [GRADED FACTS] status=unknown / "tidak tersedia", tulis "data tidak tersedia" —
JANGAN mengisi dengan angka standar.
- Jika [GRADED FACTS] menyatakan direction/status, ULANGI arah itu — jangan dibalik.
- SATUAN SETIAP ANGKA TERTULIS DI AKHIR NAMA FIELD: _gram, _persen, _rasio, _karung, _ekor, _kg,
_hari, _indeksTanpaSatuan. Pakai satuan itu persis.
- Field yang berisi "tidak tersedia" atau null memang tidak ada datanya. Tulis "data tidak tersedia"
dan JANGAN mengarang angkanya.
- Status topik tanpa data = unknown, bukan ok.
- Panen ≠ kematian; jangan hitung mortalitas dari selisih populasi awal − kini.
- DILARANG menghitung sendiri, memperkirakan, atau membalik arah tren.
FORMAT OUTPUT (JSON SAJA):
{"kesimpulan":"...","insight":"..."}
ATURAN FORMAT NARASI (WAJIB — BERLAKU UNTUK SEMUA TOPIK):
- "kesimpulan" dan "insight" adalah paragraf naratif utuh yang mengalir,
sama seperti insight halaman lain — bukan daftar, bukan poin-poin.
- DILARANG menulis baris per-metrik seperti "fcr: ..." / "bw: ..." dan
DILARANG menyalin pesan [GRADED FACTS] mentah — rangkai menjadi analisis utuh.
"""
END_CYCLE_OUTPUT_RULES = """
FORMAT OUTPUT AKHIR SIKLUS (JSON SAJA — WAJIB):
{
"kesimpulan":"...",
"masalah":["..."],
"akar_penyebab":[{"id":"...","hipotesis":"...","bukti":["..."],"confidence":"high|med|low"}],
"perbaikan_siklus_berikutnya":[{"id":"...","fase":"brooding|growth|finisher","aksi":"...","metrik_pantau":"..."}],
"insight":"..."
}
ATURAN AKHIR SIKLUS:
- Masalah = ringkas isu dari [GRADED FACTS] (critical/warning saja).
- Akar penyebab HANYA dari [ROOT CAUSE CANDIDATES] — boleh rapikan bahasa/bukti, JANGAN menambah hipotesis baru.
- Perbaikan HANYA dari [NEXT CYCLE ACTIONS] — boleh rapikan bahasa, JANGAN menambah aksi baru.
- Sertakan field "id" dari kandidat/aksi yang dipakai.
- Jika kandidat kosong dan KPI sehat, tulis kesimpulan positif singkat; masalah/akar/perbaikan boleh [].
"""
def build_insight_user_prompt(
*,
topic: str,
kandang_name: str,
report_type: str,
report_period: str,
context: dict[str, Any],
) -> str:
"""User prompt for Ollama — name as-is, never prefix 'kandang' or expose ids."""
# Drop internal ids from the JSON the model sees (defense against id=N narration).
safe_ctx = {
k: v
for k, v in context.items()
if k not in ("kandangId", "kandang_id", "cycleId", "cycle_id")
}
return (
f"Buat insight topik `{topic}` untuk unit bernama `{kandang_name}` "
f"periode `{report_type}/{report_period}`.\n"
f"NAMA WAJIB: tulis PERSIS `{kandang_name}` — jangan menambah kata 'Kandang' di depan "
f"(salah: 'Kandang {kandang_name}'; benar: '{kandang_name}'), dan jangan menyebut id.\n"
f"[Data Halaman (JSON)]\n{json.dumps(safe_ctx, ensure_ascii=False, default=str)}"
)
def strip_headerless_tables(chunk: str) -> tuple[str, int]:
lines = str(chunk or "").split("\n")
removed = 0
kept: list[str] = []
for line in lines:
trimmed = line.strip()
if trimmed and _NUMERIC_PIPE_ROW.match(trimmed):
removed += 1
continue
kept.append(line)
if removed == 0:
return chunk, 0
text = (
"\n".join(kept).strip()
+ "\n[Tabel angka tanpa judul kolom dihapus dari kutipan ini karena tidak dapat "
"dibaca dengan benar. Gunakan blok STANDAR CP 707 / GRADED FACTS untuk angka standar.]"
)
return text, removed
def _day_age(context: dict[str, Any]) -> int | None:
raw = (
context.get("hari_ke")
or context.get("hari_terakhir")
or context.get("currentDay")
or context.get("dayAge")
)
try:
day = int(raw)
except (TypeError, ValueError):
return None
return day if day > 0 else None
def prune_context_by_period(context: dict[str, Any], report_type: str, report_period: str) -> dict[str, Any]:
"""Light prune for history arrays; keep scalars. Full FE scope happens client-side."""
out = dict(context)
# Placeholder zeros (esp. weight/FCR sensors) → null so LLM sees "unavailable".
for key in (
"bobot_avg_gram",
"averageWeight",
"bobot_rata_rata_gram",
"bobot_iot_gram",
"bobot_iot_gram_terakhir",
"fcr_terakhir",
"fcr",
"eef_terakhir",
):
val = out.get(key)
if isinstance(val, (int, float)) and val <= 0:
out[key] = None
if report_type == "end_cycle":
# Deterministic weekly KPI rollup when weekly_summaries absent.
if "weekly_summaries" not in out or not out.get("weekly_summaries"):
out["weekly_summaries"] = _build_weekly_summaries(out)
if "phase_summaries" not in out or not out.get("phase_summaries"):
out["phase_summaries"] = _build_phase_summaries(out.get("weekly_summaries") or [])
out["catatan_end_cycle"] = (
"Ringkasan KPI per minggu dihitung di backend. Narasikan dari weekly_summaries + "
"phase_summaries + graded_facts + root-cause candidates; jangan menghitung ulang."
)
out["report_type"] = report_type
out["report_period"] = report_period
return out
def _row_day(row: dict[str, Any]) -> int | None:
day = row.get("hari") or row.get("day") or row.get("age")
try:
day_i = int(day)
except (TypeError, ValueError):
return None
return day_i if day_i > 0 else None
def _row_num(row: dict[str, Any], *keys: str) -> float | None:
for key in keys:
val = row.get(key)
if isinstance(val, (int, float)):
return float(val)
return None
def _build_weekly_summaries(context: dict[str, Any]) -> list[dict[str, Any]]:
histories: list[Any] = []
for key in (
"kpiSeries",
"tren_kpi_harian",
"tren_fcr_harian",
"tren_harian",
"tren_7_hari_terakhir",
"history",
):
val = context.get(key)
if isinstance(val, list) and val:
histories = val
break
by_week: dict[int, list[dict[str, Any]]] = {}
for row in histories:
if not isinstance(row, dict):
continue
day_i = _row_day(row)
if day_i is None:
continue
week = max(1, (day_i + 6) // 7)
by_week.setdefault(week, []).append(row)
summaries = []
prev_mort: float | None = None
for week, rows in sorted(by_week.items()):
rows_sorted = sorted(rows, key=lambda r: _row_day(r) or 0)
fcrs = [
v
for r in rows_sorted
if (v := _row_num(r, "fcr_aktual", "fcr", "fcr_terakhir")) is not None
]
bws = [
v
for r in rows_sorted
if (v := _row_num(r, "bobot_avg_gram", "actual", "average_weight", "bw")) is not None
]
feeds = [
v
for r in rows_sorted
if (v := _row_num(r, "pakan_kumulatif_karung", "feed_total", "karung")) is not None
]
morts = [
v
for r in rows_sorted
if (v := _row_num(r, "mortalitas_kumulatif_persen", "mortality_pct")) is not None
]
mort_end = morts[-1] if morts else None
mort_delta = None
if mort_end is not None:
mort_delta = round(mort_end - (prev_mort or 0.0), 3)
prev_mort = mort_end
cold = sum(
_row_num(r, "jam_suhu_di_bawah_standar", "temp_below_hours", "hours_temp_below") or 0
for r in rows_sorted
)
humid = sum(
_row_num(
r, "jam_kelembapan_di_atas_standar", "humidity_above_hours", "hours_humidity_above"
)
or 0
for r in rows_sorted
)
ammonia = sum(
_row_num(r, "jam_amonia_tinggi", "ammonia_high_hours", "hours_ammonia_high") or 0
for r in rows_sorted
)
summaries.append(
{
"minggu_ke": week,
"jumlah_hari_data": len(rows_sorted),
"fcr_rata_rasio": round(sum(fcrs) / len(fcrs), 3) if fcrs else None,
"fcr_akhir_rasio": fcrs[-1] if fcrs else None,
"bobot_akhir_gram": bws[-1] if bws else None,
"pakan_akhir_karung": feeds[-1] if feeds else None,
"mortalitas_kumulatif_persen": mort_end,
"mortalitas_delta_persen": mort_delta,
"jam_suhu_di_bawah_standar": round(cold, 1) if cold else 0,
"jam_kelembapan_di_atas_standar": round(humid, 1) if humid else 0,
"jam_amonia_tinggi": round(ammonia, 1) if ammonia else 0,
}
)
return summaries
def _build_phase_summaries(weekly: list[dict[str, Any]]) -> dict[str, Any]:
"""Roll weekly rows into brooding / growth / finisher buckets."""
buckets: dict[str, list[dict[str, Any]]] = {
root_cause.PHASE_BROODING: [],
root_cause.PHASE_GROWTH: [],
root_cause.PHASE_FINISHER: [],
}
for row in weekly:
if not isinstance(row, dict):
continue
week = row.get("minggu_ke")
try:
week_i = int(week)
except (TypeError, ValueError):
continue
# week 1-2 ≈ days 1-14, week 3-4 ≈ 15-28, week 5+ ≈ 29+
if week_i <= 2:
buckets[root_cause.PHASE_BROODING].append(row)
elif week_i <= 4:
buckets[root_cause.PHASE_GROWTH].append(row)
else:
buckets[root_cause.PHASE_FINISHER].append(row)
out: dict[str, Any] = {}
for phase, rows in buckets.items():
if not rows:
continue
mort_deltas = [
float(r["mortalitas_delta_persen"])
for r in rows
if isinstance(r.get("mortalitas_delta_persen"), (int, float))
]
fcrs = [
float(r["fcr_akhir_rasio"])
for r in rows
if isinstance(r.get("fcr_akhir_rasio"), (int, float))
]
out[phase] = {
"jumlah_minggu": len(rows),
"mortalitas_delta_persen": round(sum(mort_deltas), 3) if mort_deltas else None,
"fcr_akhir_rasio": fcrs[-1] if fcrs else None,
}
return out
def grade_context(context: dict[str, Any]) -> dict[str, Any]:
analysis = cp707.analyze_with_cp707_standards(context)
graded: dict[str, Any] = {
"kandangName": context.get("kandangName") or context.get("kandang_name"),
"hari_ke": analysis.get("dayAge") or _day_age(context),
"analyses": analysis.get("analyses") or {},
}
# Surface daily mortality headcount so the model cannot confuse the 5%
# cumulative threshold with "mortalitas hari ini".
if "mortalitas_hari_ini_ekor" in context:
graded["mortalitas_hari_ini_ekor"] = context.get("mortalitas_hari_ini_ekor")
# Normalize analyzer outputs into explicit direction blocks for the prompt.
for key, block in list(graded["analyses"].items()):
if not isinstance(block, dict):
continue
if "direction" not in block and block.get("deviation") is not None:
try:
dev = float(block["deviation"])
if abs(dev) <= 5:
block["direction"] = "sesuai_standar"
elif key == "fcr":
block["direction"] = "di_atas_standar" if dev > 0 else "di_bawah_standar"
else:
block["direction"] = "di_atas_standar" if dev > 0 else "di_bawah_standar"
except (TypeError, ValueError):
pass
return graded
def fetch_rag_chunks(query: str, topic: str, n_results: int = 4) -> tuple[list[str], list[dict[str, Any]]]:
base = getattr(settings, "RAG_SERVICE_URL", "") or ""
if not base:
return [], []
url = f"{base.rstrip('/')}/query"
try:
with httpx.Client(timeout=8.0) as client:
resp = client.post(
url,
json={"query": query, "topic": topic or "", "n_results": n_results, "tipe": "prosa"},
)
if resp.status_code >= 400:
logger.warning("RAG query HTTP %s", resp.status_code)
return [], []
data = resp.json()
chunks = data.get("chunks") or []
metas = data.get("metadatas") or []
sources = data.get("sources") or []
citations = []
cleaned = []
for i, chunk in enumerate(chunks):
text, _ = strip_headerless_tables(chunk)
if text.strip():
cleaned.append(text)
meta = metas[i] if i < len(metas) else {}
citations.append(
{
"source": (meta or {}).get("source") or (sources[i] if i < len(sources) else "cp707"),
"bab": (meta or {}).get("bab") or "",
"chunk_id": (meta or {}).get("chunk_index"),
"tipe": (meta or {}).get("tipe") or "prosa",
"excerpt": text[:240],
}
)
return cleaned, citations
except Exception as exc: # noqa: BLE001
logger.warning("RAG unavailable: %s", exc)
return [], []
def call_ollama(system_prompt: str, user_prompt: str) -> str | None:
model = getattr(settings, "LLM_MODEL_NAME", "qwen2.5:3b") or "qwen2.5:3b"
base = getattr(settings, "OLLAMA_BASE_URL", "http://127.0.0.1:11434") or "http://127.0.0.1:11434"
url = f"{base.rstrip('/')}/api/chat"
timeout = float(getattr(settings, "LLM_TIMEOUT_SECONDS", 1200) or 1200)
payload = {
"model": model,
"stream": False,
"options": {"temperature": 0.2},
"messages": [
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt},
],
}
try:
with httpx.Client(timeout=timeout) as client:
resp = client.post(url, json=payload)
if resp.status_code >= 400:
logger.error("Ollama HTTP %s: %s", resp.status_code, resp.text[:500])
return None
data = resp.json()
message = data.get("message") or {}
return message.get("content") or data.get("response")
except Exception as exc: # noqa: BLE001
logger.error("Ollama call failed: %s", exc)
return None
def _normalize_string_newlines(fragment: str) -> str:
"""Replace raw newlines/tabs inside JSON string literals with spaces.
qwen2.5:3b wraps long string values across real line breaks, which is
invalid JSON; without this every long narration fell back to one-liners.
"""
out: list[str] = []
in_str = False
i = 0
while i < len(fragment):
ch = fragment[i]
if in_str:
if ch == "\\" and i + 1 < len(fragment):
out.append(ch)
out.append(fragment[i + 1])
i += 2
continue
if ch == '"':
in_str = False
out.append(ch)
i += 1
continue
if ch in "\n\r\t":
out.append(" ")
i += 1
continue
elif ch == '"':
in_str = True
out.append(ch)
i += 1
return "".join(out)
def parse_llm_json(text: str) -> dict[str, str] | None:
if not text:
return None
# Strip fence MARKERS only (FE insightParse parity) — removing the whole
# ```…``` block deleted the JSON itself and forced the local fallback.
cleaned = re.sub(r"^```(?:json)?\s*", "", text.strip(), flags=re.I)
cleaned = re.sub(r"\s*```\s*$", "", cleaned)
# Tolerate small-model JSON noise (FE parser parity): // and /* */ comments,
# trailing commas. Lookbehind keeps "http://" URLs intact.
cleaned = re.sub(r"(?<!:)//[^\n]*", "", cleaned)
cleaned = re.sub(r"/\*.*?\*/", "", cleaned, flags=re.S)
cleaned = re.sub(r",(\s*[}\]])", r"\1", cleaned)
first = cleaned.find("{")
last = cleaned.rfind("}")
if first < 0 or last <= first:
# No JSON object — the model wrote free prose; keep it as the narrative
# instead of dropping it to the per-metric local fallback.
prose = cleaned.strip()
if len(prose) < 100:
return None
parts = re.split(r"(?<=[.!?])\s+", prose, maxsplit=1)
kesimpulan = parts[0].strip()
insight = parts[1].strip() if len(parts) > 1 else kesimpulan
return {"kesimpulan": kesimpulan, "insight": insight}
fragment = cleaned[first : last + 1]
try:
parsed = json.loads(fragment)
except json.JSONDecodeError:
try:
parsed = json.loads(_normalize_string_newlines(fragment))
except json.JSONDecodeError:
return None
if not isinstance(parsed, dict):
return None
kesimpulan = (
parsed.get("kesimpulan")
or parsed.get("ringkasan")
or parsed.get("summary")
or ""
)
insight = (
parsed.get("insight")
or parsed.get("rekomendasi")
or parsed.get("insights")
or ""
)
if not kesimpulan and not insight:
# Model invented its own keys (e.g. per-metric fcr/bw/mortality) — fold
# the values into narrative paragraphs rather than failing to fallback.
values = [str(v).strip() for v in parsed.values() if isinstance(v, str) and v.strip()]
if not values:
return None
return {
"kesimpulan": values[0],
"insight": "\n\n".join(values[1:]) or values[0],
}
return {
"kesimpulan": str(kesimpulan) or "Model tidak mengembalikan kesimpulan eksplisit.",
"insight": str(insight) or "Model tidak mengembalikan rekomendasi eksplisit.",
}
def parse_end_cycle_llm_json(text: str) -> dict[str, Any] | None:
"""Parse end-cycle structured JSON; tolerate partial fields."""
if not text:
return None
cleaned = re.sub(r"<think>[\s\S]*?</think>", "", text, flags=re.I)
cleaned = re.sub(r"^```(?:json)?\s*", "", cleaned.strip(), flags=re.I)
cleaned = re.sub(r"\s*```\s*$", "", cleaned)
first = cleaned.find("{")
last = cleaned.rfind("}")
if first < 0 or last <= first:
return None
try:
parsed = json.loads(cleaned[first : last + 1])
except json.JSONDecodeError:
return None
if not isinstance(parsed, dict):
return None
# Accept if any of the end-cycle keys or daily keys exist.
keys = (
"kesimpulan",
"ringkasan",
"summary",
"masalah",
"akar_penyebab",
"perbaikan_siklus_berikutnya",
"insight",
)
if not any(k in parsed for k in keys):
return None
return parsed
def local_fallback_insight(graded: dict[str, Any], topic: str) -> dict[str, str]:
analyses = graded.get("analyses") or {}
lines = []
for name, block in analyses.items():
if isinstance(block, dict) and block.get("message"):
lines.append(f"{name}: {block['message']}")
if not lines:
return {
"kesimpulan": f"Data untuk topik {topic} tidak cukup untuk dianalisis (status unknown).",
"insight": "Lengkapi data operasional kandang, lalu generate ulang. Jangan mengarang angka.",
}
return {
"kesimpulan": lines[0],
"insight": "\n".join(lines[1:]) if len(lines) > 1 else lines[0],
}
_SOFT_MORTALITY_WORD = re.compile(
r"\b(rendah|normal|baik|aman|terkendali|sedikit)\b",
re.IGNORECASE,
)
_MENTIONS_MORTALITY = re.compile(r"mortalitas", re.IGNORECASE)
def _text_soft_labels_mortality(text: str) -> bool:
"""True when narrative uses soft adjectives in a mortality sentence."""
if not text or not _MENTIONS_MORTALITY.search(text):
return False
# Prefer soft word near "mortalitas" (same sentence-ish window).
for match in _MENTIONS_MORTALITY.finditer(text):
start = max(0, match.start() - 60)
end = min(len(text), match.end() + 80)
if _SOFT_MORTALITY_WORD.search(text[start:end]):
return True
return bool(_SOFT_MORTALITY_WORD.search(text) and _MENTIONS_MORTALITY.search(text))
def enforce_mortality_wording(
parsed: dict[str, Any],
graded: dict[str, Any] | None,
) -> dict[str, Any]:
"""Rewrite LLM text that calls high mortality 'rendah/normal' (small-model slip)."""
mort = ((graded or {}).get("analyses") or {}).get("mortality") or {}
status = str(mort.get("status") or "").lower()
if status not in ("warning", "critical"):
return parsed
label = mort.get("severity_label") or ("SANGAT TINGGI" if status == "critical" else "TINGGI")
actual = mort.get("actual_pct")
pct = f"{float(actual):.2f}%" if isinstance(actual, (int, float)) else None
message = str(mort.get("message") or "").strip()
out = dict(parsed)
kesimpulan = str(out.get("kesimpulan") or "")
insight = str(out.get("insight") or "")
if not (_text_soft_labels_mortality(kesimpulan) or _text_soft_labels_mortality(insight)):
return out
# Prefer the deterministic graded message; keep non-mortality insight if clean.
if message:
out["kesimpulan"] = message
else:
pct_part = f" ({pct})" if pct else ""
out["kesimpulan"] = (
f"Mortalitas kumulatif{pct_part} adalah {label} menurut ambang CP 707."
)
if _text_soft_labels_mortality(insight):
out["insight"] = (
"Segera tinjau penyebab kematian (nekropsi bila kritis), perketat biosekuriti, "
"dan pantau mortalitas harian hingga tren menurun."
)
return out
_DUP_KANDANG = re.compile(r"\b[Kk]andang(?:\s+[Kk]andang)+\b")
# Line-start metric labels ("fcr: ...", "bw: ...") — per-metric one-liners the
# model (or local fallback) may emit instead of narrative paragraphs.
_METRIC_LABEL_LINE = re.compile(
r"(?mi)^\s*(?:fcr|eef|bw|bobot|mortality|mortalitas|environment|ip)\s*:\s*(?=\S)"
)
def collapse_duplicate_kandang(text: str) -> str:
"""Fix LLM slip 'Kandang Kandang 2' → 'Kandang 2' (name already includes Kandang)."""
if not text:
return text
return _DUP_KANDANG.sub(lambda m: "Kandang" if m.group(0)[0].isupper() else "kandang", text)
def strip_metric_label_lines(text: str) -> str:
"""Drop `metric: ` line labels and flow the lines into one paragraph."""
if not text:
return text
stripped = _METRIC_LABEL_LINE.sub("", text)
if stripped == text:
return text
return re.sub(r"\s*\n+\s*", " ", stripped).strip()
def sanitize_insight_narrative(parsed: dict[str, Any]) -> dict[str, Any]:
"""Deterministic cleanup of narrative fields after the LLM."""
out = dict(parsed)
for key in ("kesimpulan", "insight", "summary", "insight_text"):
if key in out and out[key] is not None:
text = collapse_duplicate_kandang(str(out[key]))
out[key] = strip_metric_label_lines(text)
return out
def lookup_cached(
*,
cycle_id: int,
kandang_id: int,
topic: str,
report_type: str,
report_period: str,
) -> AIInsight | None:
return (
AIInsight.objects.filter(
cycle_id=cycle_id,
kandang_id=kandang_id,
topic=topic,
report_type=report_type,
report_period=report_period or "current",
)
.order_by("-updated_at")
.first()
)
def get_cached_insight(
*,
cycle_id: int,
kandang_id: int,
topic: str,
report_type: str = "page",
report_period: str = "current",
) -> dict[str, Any] | None:
row = (
AIInsight.objects.filter(
cycle_id=cycle_id,
kandang_id=kandang_id,
topic=topic,
report_type=report_type,
report_period=report_period or "current",
)
.order_by("-updated_at")
.first()
)
if row is None:
return None
return serialize_insight(row)
def encode_citations(citations: list[Any] | None) -> str:
"""Persist citations as a JSON array string in TEXT column."""
if not citations:
return "[]"
return json.dumps(citations, ensure_ascii=False, default=str)
def decode_citations(raw: Any) -> list[dict[str, Any]]:
"""Load citations from TEXT (JSON string) or legacy list."""
if raw is None or raw == "":
return []
if isinstance(raw, list):
data = raw
elif isinstance(raw, str):
try:
data = json.loads(raw)
except json.JSONDecodeError:
return []
else:
return []
if not isinstance(data, list):
return []
# Normalize citation keys for FE (chapter alias).
norm_citations: list[dict[str, Any]] = []
for c in data:
if not isinstance(c, dict):
continue
item = dict(c)
if "chapter" not in item and item.get("bab"):
item["chapter"] = item["bab"]
norm_citations.append(item)
return norm_citations
def serialize_insight(row: AIInsight) -> dict[str, Any]:
norm_citations = decode_citations(row.citations)
summary = collapse_duplicate_kandang(row.summary or "")
insight_body = collapse_duplicate_kandang(row.insight_text or "")
insight_text = "\n\n".join(
part for part in [summary, insight_body] if part
) or insight_body or ""
structured = root_cause.parse_embedded_end_cycle(insight_body)
if structured and summary and not structured.get("kesimpulan"):
structured = {**structured, "kesimpulan": summary}
elif structured:
structured = sanitize_insight_narrative(structured)
return {
"success": True,
"id": row.pk,
"cycle": row.cycle_id,
"kandang": row.kandang_id,
"topic": row.topic,
"report_type": row.report_type,
"report_period": row.report_period,
"summary": summary,
"insight": insight_body,
"insight_text": insight_text,
"structured_end_cycle": structured,
"alert": row.alert or ALERT_UNKNOWN,
"citations": norm_citations,
"date": row.date.isoformat() if row.date else None,
"created_at": row.created_at.isoformat() if row.created_at else None,
"updated_at": row.updated_at.isoformat() if row.updated_at else None,
}
def generate_insight(
*,
cycle_id: int,
kandang_id: int,
topic: str,
context: dict[str, Any],
report_type: str = "page",
report_period: str = "current",
force_refresh: bool = False,
) -> dict[str, Any]:
if not kandang_id:
raise ValueError("kandang_id wajib — insight tidak boleh untuk semua kandang")
if not cycle_id:
raise ValueError("cycle_id wajib")
if not topic:
raise ValueError("topic wajib")
try:
cycle = Cycle.objects.select_related("kandang").get(pk=cycle_id)
except Cycle.DoesNotExist as exc:
raise ValueError("Cycle not found") from exc
try:
kandang = Kandang.objects.get(pk=kandang_id)
except Kandang.DoesNotExist as exc:
raise ValueError("Kandang not found") from exc
if cycle.kandang_id != kandang_id:
raise ValueError("kandang_id tidak cocok dengan cycle")
report_period = report_period or "current"
report_type = report_type or "page"
if not force_refresh:
cached = get_cached_insight(
cycle_id=cycle_id,
kandang_id=kandang_id,
topic=topic,
report_type=report_type,
report_period=report_period,
)
if cached:
return cached
ctx = dict(context or {})
ctx.setdefault("kandangId", kandang_id)
ctx.setdefault("kandang_id", kandang_id)
ctx.setdefault("kandangName", kandang.kandang_name)
ctx.setdefault("cycleId", cycle_id)
ctx = prune_context_by_period(ctx, report_type, report_period)
graded = grade_context(ctx)
day = graded.get("hari_ke")
standard_block = cp707.build_cp707_standard_block(ctx) or cp707.build_book_reference_context(day or 1)
is_end_cycle = report_type == "end_cycle"
hypotheses: list[dict[str, Any]] = []
next_actions: list[dict[str, Any]] = []
masalah: list[str] = []
if is_end_cycle:
hypotheses = root_cause.build_root_cause_hypotheses(graded, ctx)
next_actions = root_cause.build_next_cycle_actions(hypotheses)
masalah = root_cause.build_masalah_from_graded(graded)
# Prefer phase-targeted RAG for end-cycle remediation.
top_fase = next((h.get("fase") for h in hypotheses if h.get("fase")), None)
rag_query = (
f"panduan CP 707 perbaikan {top_fase or 'broiler'} siklus berikutnya "
f"umur hari ke-{day or '?'}"
)
else:
rag_query = f"panduan manajemen broiler CP 707 untuk {topic} umur hari ke-{day or '?'}"
chunks, citations = fetch_rag_chunks(rag_query, topic)
system_prompt = (
"Anda adalah asisten farm broiler on-premise. Tugas Anda HANYA menulis narasi "
"dari GRADED FACTS + standar CP 707 + cuplikan SOP. Jangan menghitung ulang.\n"
f"{ANTI_HALLUCINATION_RULES}\n"
)
if is_end_cycle:
# End-cycle replaces the 2-field daily format instructions.
system_prompt = system_prompt.replace(
'FORMAT OUTPUT (JSON SAJA):\n{"kesimpulan":"...","insight":"..."}',
END_CYCLE_OUTPUT_RULES.strip(),
)
system_prompt += (
f"{standard_block}\n"
f"[GRADED FACTS]\n{json.dumps(graded, ensure_ascii=False, default=str)}\n"
)
if is_end_cycle:
system_prompt += (
f"\n[ROOT CAUSE CANDIDATES]\n"
f"{json.dumps(hypotheses, ensure_ascii=False, default=str)}\n"
f"\n[NEXT CYCLE ACTIONS]\n"
f"{json.dumps(next_actions, ensure_ascii=False, default=str)}\n"
)
if chunks:
system_prompt += "\n[CUPLIKAN SOP CP 707 — prosa]\n" + "\n---\n".join(chunks[:4])
user_prompt = build_insight_user_prompt(
topic=topic,
kandang_name=kandang.kandang_name,
report_type=report_type,
report_period=report_period,
context=ctx,
)
raw = call_ollama(system_prompt, user_prompt)
if is_end_cycle:
parsed_raw = parse_end_cycle_llm_json(raw or "")
if parsed_raw:
# Apply mortality wording on narrative fields before sanitize.
soft = {
"kesimpulan": str(parsed_raw.get("kesimpulan") or ""),
"insight": str(parsed_raw.get("insight") or ""),
}
soft = enforce_mortality_wording(soft, graded)
soft = sanitize_insight_narrative(soft)
parsed_raw["kesimpulan"] = soft.get("kesimpulan")
parsed_raw["insight"] = soft.get("insight")
payload = root_cause.sanitize_end_cycle_payload(
parsed_raw,
hypotheses=hypotheses,
actions=next_actions,
masalah=masalah,
)
payload = sanitize_insight_narrative(payload)
else:
payload = root_cause.local_fallback_end_cycle(graded, ctx)
summary = collapse_duplicate_kandang(str(payload.get("kesimpulan") or ""))
insight_text = collapse_duplicate_kandang(
root_cause.format_end_cycle_insight_text(payload)
)
else:
parsed = parse_llm_json(raw or "")
if not parsed:
logger.warning(
"Daily LLM output unparseable (topic=%s); local fallback. raw[:400]=%r",
topic,
(raw or "")[:400],
)
parsed = local_fallback_insight(graded, topic)
else:
parsed = enforce_mortality_wording(parsed, graded)
parsed = sanitize_insight_narrative(parsed)
insight_text = parsed["insight"]
summary = parsed["kesimpulan"]
alert = alert_from_graded(graded)
row, _created = AIInsight.objects.update_or_create(
cycle=cycle,
kandang=kandang,
topic=topic,
report_type=report_type,
report_period=report_period,
defaults={
"date": timezone.localdate(),
"insight_text": insight_text,
"summary": summary,
"alert": alert,
"citations": encode_citations(citations),
},
)
return serialize_insight(row)