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

1359 lines
54 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
import threading
import time
import uuid
from pathlib import Path
from typing import Any, Callable
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__)
# Plan §5: one LLM inference at a time (CPU server; prevents thread thrash).
_INFERENCE_LOCK = threading.Lock()
_NUMERIC_PIPE_ROW = re.compile(r"^[\d.,]+(\s*\|\s*[\d.,]*)+$")
DAILY_JSON_MARKER = "<!--daily_json-->"
def parse_embedded_daily(text: str) -> dict[str, Any] | None:
"""Extract structured daily payload (akar_masalah, insight) from stored insight_text."""
if not text or DAILY_JSON_MARKER not in text:
return None
raw = text.split(DAILY_JSON_MARKER, 1)[1].strip()
try:
data = json.loads(raw)
except json.JSONDecodeError:
return None
return data if isinstance(data, dict) else None
# 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.
ATURAN STANDAR CP 707 — DILARANG MENGARANG (WAJIB DIPATUHI):
- Blok STANDAR CP 707 menyatakan PERSIS angka mana yang ada di buku dan mana yang TIDAK.
- Jika standar menyatakan "TIDAK tercantum di buku" atau "N/A" — DILARANG KERAS mengarang
angka standar sendiri (mis. "standar mortalitas 9.32%" atau "standar EEF 235").
- EEF/IP standar HANYA ada hari ke-7 s/d ke-37 (rentang 327–380). Di luar rentang itu
standar TIDAK ADA — jangan mengarang.
- Mortalitas kumulatif standar hanya sampai hari ke-37 (4.65%). Setelah hari 37 buku hanya
memperkirakan +0.1%/hari — ini BUKAN standar resmi. Jika hari >37, tulis "standar tidak
tercantum di buku untuk hari ini" dan bandingkan dengan 4.65% + estimasi, bukan mengarang.
- FCR standar hari >37 gunakan fallback 1.65 (bukan standar buku) — sebutkan "fallback" bukan standar.
- Bobot standar hanya hari 1–37. Di luar itu "tidak tercantum di buku".
FORMAT OUTPUT (JSON SAJA):
{"kesimpulan":"...","akar_masalah":"...","insight":"..."}
"""
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],
graded: dict[str, Any] | None = None,
is_end_cycle: bool = False,
hypotheses: list[dict[str, Any]] | None = None,
) -> str:
"""Build ordered user prompt template for Ollama.
Template order (recency = last wins):
1. Header: topic, kandang name, period, naming rules
2. Data Halaman (JSON) — page data only, no IDs
3. STATUS AKTUAL (alert from graded) + FAKTA GRADED (status lines, no unknown when known exists)
4. Root Cause / ANALISIS AKAR MASALAH (end_cycle only, from hypotheses)
5. Length contract + field role contract
6. Output schema (JSON only, placeholders as <...>)
"""
# 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")
}
lines: list[str] = []
# 1. Header
lines.append(
f"Buat insight topik `{topic}` untuk unit bernama `{kandang_name}` "
f"periode `{report_type}/{report_period}`."
)
lines.append(
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."
)
# 2. Data Halaman (JSON)
lines.append(f"[Data Halaman (JSON)]\n{json.dumps(safe_ctx, ensure_ascii=False, default=str)}")
# 3. STATUS AKTUAL + FAKTA GRADED
if graded:
analyses = graded.get("analyses") or {}
known_lines = [
f"{name} = {block.get('status')}: {block.get('message')}"
for name, block in analyses.items()
if isinstance(block, dict)
and block.get("message")
and block.get("status") != "unknown"
]
unknown_lines = [
f"{name} = {block.get('status')}: {block.get('message')}"
for name, block in analyses.items()
if isinstance(block, dict)
and block.get("message")
and block.get("status") == "unknown"
]
status_lines = known_lines or unknown_lines
actual_alert = alert_from_graded(graded)
lines.append(
"\nKONTEKS: semua data adalah ternak AYAM BROILER (unggas) di kandang — "
"BUKAN tanaman/pertanian. 'Panen' = panen ayam. 'Bobot' = gram per ekor "
"(bukan per karung). 'ADG' = kenaikan bobot harian ayam."
)
lines.append(
f"STATUS AKTUAL: {actual_alert}. Kesimpulan dan insight HARUS konsisten "
"dengan status ini: bila warning/critical, sebut masalahnya dan bandingkan "
"dengan angka standar CP 707; dilarang menyebut ideal/sehat/baik/aman "
"bila STATUS AKTUAL bukan healthy."
)
if status_lines:
lines.append("FAKTA GRADED:\n- " + "\n- ".join(status_lines))
if topic == "hitung_karung":
lines.append(
"\nFOKUS KHUSUS HITUNG KARUNG: Narasi dan kesimpulan WAJIB fokus HANYA pada stok/saldo karung pakan, "
"karung masuk, saldo awal, dan pakan dituang. JANGAN menyebut, mencari, atau mengeluhkan ketiadaan data mortalitas, "
"bobot ayam, atau FCR karena metrik tersebut bukan bagian dari topik Hitung Karung."
)
# 4. Root Cause / ANALISIS AKAR MASALAH (end_cycle only)
if is_end_cycle and hypotheses:
lines.append("\n[ANALISIS AKAR MASALAH]")
for h in hypotheses:
hid = h.get("id", "")
label = h.get("hipotesis", "")
bukti = h.get("bukti", [])
conf = h.get("confidence", "")
fase = h.get("fase", "")
line = f"- {label} ({conf}, fase: {fase})" if conf or fase else f"- {label}"
lines.append(line)
for b in bukti:
lines.append(f" - {b}")
# 5. Length contract + field role contract
actual_alert = alert_from_graded(graded) if graded else ALERT_UNKNOWN
if is_end_cycle:
lines.append(
"\nPANJANG WAJIB: kesimpulan = 1-3 kalimat PENILAIAN DATA SAJA (tanpa sebab/aksi). "
"masalah = bullet isu dari FAKTA GRADED. akar_penyebab = dari ANALISIS AKAR MASALAH saja. "
"perbaikan_siklus_berikutnya = dari aksi yang tersedia saja. "
"insight = narasi 3-6 kalimat berurutan: "
"(1) angka aktual vs standar CP 707, (2) penyebab/implikasi, (3) tindakan konkret. "
"Satu kalimat singkat = gagal. DILARANG menulis kata aksi (lakukan, periksa, sebaiknya, "
"konsultasikan) di kesimpulan — insight HARUS mengandung minimal satu kata aksi."
)
else:
if actual_alert in (ALERT_WARNING, ALERT_CRITICAL):
lines.append(
f"\nPANJANG WAJIB (STATUS AKTUAL ADALAH {actual_alert.upper()} — ADA ANOMALI/MASALAH): "
"kesimpulan = 1-3 kalimat PENILAIAN DATA SAJA (tanpa sebab/aksi, sebutkan angka aktual vs standar CP 707). "
"akar_masalah = 1-3 kalimat ANALISIS AKAR PENYEBAB & FAKTOR PEMICU DEVIASI (soroti anomali mortalitas, fluktuasi konsumsi pakan, atau deviasi mikroklimat suhu/kelembapan). DILARANG KERAS menulis 'kondisi optimal' atau 'tidak terdeteksi anomali' karena status data sedang bermasalah. "
"insight = narasi 3-5 kalimat REKOMENDASI TINDAKAN KONKRET & PANDUAN SOP CP 707. "
"Satu kalimat singkat = gagal. DILARANG menulis kata aksi (lakukan, periksa, sebaiknya, "
"konsultasikan) di kesimpulan — insight HARUS mengandung minimal satu kata aksi."
)
elif actual_alert == ALERT_HEALTHY:
lines.append(
f"\nPANJANG WAJIB (STATUS AKTUAL ADALAH {actual_alert.upper()} — KONDISI OPTIMAL/NORMAL): "
"kesimpulan = 1-3 kalimat PENILAIAN DATA SAJA (konfirmasi performa sesuai standar CP 707, tanpa sebab/aksi). "
"akar_masalah = 1-2 kalimat konfirmasi bahwa parameter lingkungan dan pakan terkendali dengan baik tanpa indikasi deviasi. "
"insight = narasi 3-5 kalimat panduan mempertahankan SOP CP 707 dan monitoring rutin. "
"Satu kalimat singkat = gagal. DILARANG menulis kata aksi (lakukan, periksa, sebaiknya, "
"konsultasikan) di kesimpulan — insight HARUS mengandung minimal satu kata aksi."
)
else:
lines.append(
"\nPANJANG WAJIB: kesimpulan = 1-3 kalimat PENILAIAN DATA SAJA (tanpa sebab/aksi). "
"akar_masalah = 1-3 kalimat ANALISIS AKAR PENYEBAB / FAKTOR PEMICU berdasarkan data yang tersedia. "
"insight = narasi 3-5 kalimat REKOMENDASI TINDAKAN KONKRET & PANDUAN SOP CP 707. "
"Satu kalimat singkat = gagal. DILARANG menulis kata aksi (lakukan, periksa, sebaiknya, "
"konsultasikan) di kesimpulan — insight HARUS mengandung minimal satu kata aksi."
)
# 6. Output schema
if is_end_cycle:
lines.append(
'\nBALAS HANYA JSON persis: '
'{"kesimpulan":"<ringkasan penilaian>",'
'"masalah":["<isu dari FAKTA GRADED>"],'
'"akar_penyebab":[{"id":"<id dari ANALISIS AKAR MASALAH>","hipotesis":"<label>","bukti":["..."],"confidence":"high|med|low","fase":"brooding|growth|finisher"}],'
'"perbaikan_siklus_berikutnya":[{"id":"<id aksi>","fase":"brooding|growth|finisher","aksi":"...","metrik_pantau":"..."}],'
'"insight":"<narasi panjang 3-6 kalimat>"}'
' - JANGAN salin teks <> dari contoh; tulis kalimat lengkap sendiri.'
)
else:
if actual_alert in (ALERT_WARNING, ALERT_CRITICAL):
lines.append(
'\nBALAS HANYA JSON persis: '
'{"kesimpulan":"<ringkasan penilaian data aktual vs standar 1-3 kalimat>",'
'"akar_masalah":"<analisis faktor pemicu / akar masalah deviasi dari data faktual>",'
'"insight":"<narasi rekomendasi 3-5 kalimat: tindakan konkret & panduan SOP CP 707>"}'
' - JANGAN salin teks <> dari contoh; tulis kalimat lengkap sendiri.'
)
elif actual_alert == ALERT_HEALTHY:
lines.append(
'\nBALAS HANYA JSON persis: '
'{"kesimpulan":"<ringkasan penilaian data sesuai standar CP 707 1-3 kalimat>",'
'"akar_masalah":"<konfirmasi parameter operasional terkendali baik tanpa anomali>",'
'"insight":"<narasi rekomendasi 3-5 kalimat: panduan SOP CP 707 & monitoring berkala>"}'
' - JANGAN salin teks <> dari contoh; tulis kalimat lengkap sendiri.'
)
else:
lines.append(
'\nBALAS HANYA JSON persis: '
'{"kesimpulan":"<ringkasan penilaian data 1-3 kalimat>",'
'"akar_masalah":"<analisis faktor pemicu berdasarkan data yang ada>",'
'"insight":"<narasi rekomendasi 3-5 kalimat: tindakan konkret & panduan SOP CP 707>"}'
' - JANGAN salin teks <> dari contoh; tulis kalimat lengkap sendiri.'
)
return "\n".join(lines)
def build_insight_user_prompt_legacy(
*,
topic: str,
kandang_name: str,
report_type: str,
report_period: str,
context: dict[str, Any],
) -> str:
"""Legacy user prompt (kept for backward compatibility / tests)."""
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:
for key in ("hari_ke", "hari_terakhir", "currentDay", "dayAge"):
if key in context and context[key] is not None:
raw = context[key]
try:
day = int(raw)
if day >= 0:
return day
except (TypeError, ValueError):
continue
return 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."
)
else:
# Non-end_cycle (page / daily): prune raw history lists to latest 7 days max
# so LLM prompt does not bloat into thousands of tokens.
for history_key in (
"tren_fcr_harian",
"tren_eef_harian",
"tren_populasi_harian",
"tren_harian",
"kpiSeries",
"tren_kpi_harian",
"history",
):
val = out.get(history_key)
if isinstance(val, list) and len(val) > 7:
out[history_key] = val[-7:]
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], topic: str | None = None) -> dict[str, Any]:
analysis = cp707.analyze_with_cp707_standards(context, topic)
day_val = analysis.get("dayAge") if analysis.get("dayAge") is not None else _day_age(context)
graded: dict[str, Any] = {
"kandangName": context.get("kandangName") or context.get("kandang_name"),
"hari_ke": day_val,
"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) -> tuple[str | None, dict[str, Any]]:
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,
# format:"json" (grammar) keeps the model from rambling in prose; the schema
# reminder appended AFTER the page-data JSON keeps it from copying that JSON.
"format": "json",
"options": {
"temperature": float(getattr(settings, "LLM_TEMPERATURE", 0.2) or 0.2),
"seed": int(getattr(settings, "LLM_SEED", 42) or 42),
"num_ctx": int(getattr(settings, "LLM_NUM_CTX", 8192) or 8192),
"num_predict": int(getattr(settings, "LLM_NUM_PREDICT", 768) or 768),
"num_thread": int(getattr(settings, "LLM_NUM_THREAD", 4) or 4),
},
"messages": [
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt},
],
}
try:
with httpx.Client(timeout=timeout) as client:
# Plan §5: serialize LLM calls so concurrent generates don't thrash CPU.
with _INFERENCE_LOCK:
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 {}
content = message.get("content") or data.get("response")
usage = {
"prompt_tokens": data.get("prompt_eval_count"),
"completion_tokens": data.get("eval_count"),
}
return content, usage
except Exception as exc: # noqa: BLE001
logger.error("Ollama call failed: %s", exc)
return None, {}
def parse_llm_json(text: str) -> dict[str, str] | None:
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
kesimpulan = (
parsed.get("kesimpulan")
or parsed.get("ringkasan")
or parsed.get("summary")
or ""
)
akar_masalah = (
parsed.get("akar_masalah")
or parsed.get("root_cause")
or parsed.get("akar_penyebab")
or parsed.get("penyebab")
or parsed.get("faktor_penyebab")
or parsed.get("faktor_pemicu")
or ""
)
insight = (
parsed.get("insight")
or parsed.get("rekomendasi")
or parsed.get("insights")
or parsed.get("saran")
or ""
)
if not kesimpulan and not insight and not akar_masalah:
return None
return {
"kesimpulan": str(kesimpulan) or "Model tidak mengembalikan kesimpulan eksplisit.",
"akar_masalah": str(akar_masalah).strip() if akar_masalah else "",
"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
_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))
_OPTIMAL_CONTRADICTION = re.compile(
r"\b(kondisi\s+operasional\s+optimal|tidak\s+terdeteksi\s+(?:indikasi\s+)?anomali|semua\s+parameter\s+normal|kondisi\s+optimal|operasional\s+optimal|tidak\s+ada\s+anomali)\b",
re.IGNORECASE,
)
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' or claims optimal when critical."""
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 "")
akar_masalah = str(out.get("akar_masalah") or "")
insight = str(out.get("insight") or "")
has_soft_kesimpulan = _text_soft_labels_mortality(kesimpulan)
has_soft_akar = _text_soft_labels_mortality(akar_masalah) or bool(_OPTIMAL_CONTRADICTION.search(akar_masalah))
has_soft_insight = _text_soft_labels_mortality(insight)
if not (has_soft_kesimpulan or has_soft_akar or has_soft_insight):
return out
pct_part = f" ({pct})" if pct else ""
if has_soft_kesimpulan:
if message:
out["kesimpulan"] = message
else:
out["kesimpulan"] = (
f"Mortalitas kumulatif{pct_part} adalah {label} menurut ambang CP 707."
)
if has_soft_akar:
out["akar_masalah"] = (
f"Lonjakan mortalitas kumulatif{pct_part} yang tergolong {label} menjadi pemicu deviasi performa, mengindikasikan adanya stres lingkungan, tantangan biosekuriti, atau paparan patogen."
)
if has_soft_insight:
out["insight"] = (
"Segera lakukan nekropsi pada ayam mati untuk identifikasi patogen, perketat biosekuriti kandang, "
"dan pantau mortalitas harian hingga tren menurun."
)
return out
def enforce_bw_wording(
parsed: dict[str, Any],
graded: dict[str, Any] | None,
) -> dict[str, Any]:
"""Ensure body weight status and narrative match graded facts."""
bw = ((graded or {}).get("analyses") or {}).get("bw") or {}
status = str(bw.get("status") or "").lower()
if not status:
return parsed
message = str(bw.get("message") or "").strip()
actual = bw.get("actual")
std = bw.get("standard")
out = dict(parsed)
kesimpulan = str(out.get("kesimpulan") or "")
akar_masalah = str(out.get("akar_masalah") or "")
# If graded status is OK, but summary claims it is below standard / critical:
if status == "ok" and any(
w in kesimpulan.lower()
for w in ("di bawah standar", "sangat rendah", "critical", "kritis", "jauh di bawah")
):
if message:
out["kesimpulan"] = message
elif actual is not None and std is not None:
out["kesimpulan"] = f"Bobot badan ayam ({actual:g}g) sesuai dengan target standar ({std:g}g)."
# If graded status is warning/critical for BW, but akar_masalah claims optimal:
if status in ("warning", "critical") and _OPTIMAL_CONTRADICTION.search(akar_masalah):
diff_str = f" ({actual:g}g vs {std:g}g)" if actual is not None and std is not None else ""
out["akar_masalah"] = (
f"Pertumbuhan bobot badan tertinggal dari target standar CP 707{diff_str}, "
"dipicu oleh fluktuasi asupan pakan harian atau ketidaksesuaian mikroklimat kandang."
)
return out
_DUP_KANDANG = re.compile(r"\b[Kk]andang(?:\s+[Kk]andang)+\b")
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 _placeholder_text(value: str) -> bool:
"""True when a narrative field is empty / schema placeholder (e.g. '...')."""
return len(value.strip().strip(".…<>- \t")) < 20
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", "akar_masalah", "insight", "summary", "insight_text"):
if key in out and out[key] is not None:
out[key] = collapse_duplicate_kandang(str(out[key]))
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_raw = collapse_duplicate_kandang(row.insight_text or "")
structured = root_cause.parse_embedded_end_cycle(insight_raw)
if structured and summary and not structured.get("kesimpulan"):
structured = {**structured, "kesimpulan": summary}
elif structured:
structured = sanitize_insight_narrative(structured)
embedded_daily = parse_embedded_daily(insight_raw)
akar_masalah = ""
if embedded_daily:
akar_masalah = str(embedded_daily.get("akar_masalah") or "").strip()
insight_body = str(
embedded_daily.get("insight") or insight_raw.split(DAILY_JSON_MARKER, 1)[0]
).strip()
elif structured:
insight_body = insight_raw.split(root_cause.END_CYCLE_JSON_MARKER, 1)[0].strip()
else:
insight_body = insight_raw
insight_body = collapse_duplicate_kandang(insight_body)
akar_masalah = collapse_duplicate_kandang(akar_masalah)
parts = [summary]
if akar_masalah:
parts.append(akar_masalah)
if insight_body:
parts.append(insight_body)
insight_text = "\n\n".join(part for part in parts if part) or insight_body or ""
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,
"akar_masalah": akar_masalah,
"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,
}
# Plan Stage 1: graded status -> deterministic RAG keywords (no LLM query rewrite).
_DIAGNOSTIC_KEYWORDS = {
"mortality": "penanganan mortalitas kumulatif deplesi penyebab kematian brooding",
"fcr": "target FCR pakan harian konsumsi pakan efisiensi ventilasi",
"bw": "target bobot badan ADG pertumbuhan standar mingguan",
"iot": "standar ventilasi suhu kelembapan amonia sekam basah kandang",
"environment": "standar ventilasi suhu kelembapan amonia sekam basah kandang",
"feed_balance": "manajemen stok pakan karung pakan harian konsumsi tempat pakan",
}
def diagnose_rag_keywords(graded: dict[str, Any]) -> str:
"""Stage 1: anomaly keywords from graded blocks; '' when everything ok/unknown."""
analyses = (graded or {}).get("analyses") or {}
clues: list[str] = []
for key, block in analyses.items():
if not isinstance(block, dict):
continue
if str(block.get("status") or "").lower() in ("warning", "critical"):
hint = _DIAGNOSTIC_KEYWORDS.get(key)
if hint and hint not in clues:
clues.append(hint)
return " ".join(clues)
def _notify(stage_cb: Callable[[str], None] | None, stage: str) -> None:
if not stage_cb:
return
try:
stage_cb(stage)
except Exception: # noqa: BLE001 - progress UI must never break generation
logger.debug("stage_cb(%s) failed", stage, exc_info=True)
def _log_trace(record: dict[str, Any]) -> None:
"""Plan Stage 5: append-only audit trail for replay/debug."""
try:
path = Path(settings.BASE_DIR) / "logs" / "rag_trace.jsonl"
path.parent.mkdir(parents=True, exist_ok=True)
with path.open("a", encoding="utf-8") as fh:
fh.write(json.dumps(record, ensure_ascii=False, default=str) + "\n")
except Exception: # noqa: BLE001 - logging must never fail the request
logger.warning("rag_trace write failed", exc_info=True)
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,
stage_cb: Callable[[str], None] | None = None,
) -> 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)
t_start = time.monotonic()
_notify(stage_cb, "grading")
graded = grade_context(ctx, topic)
t_after_grade = time.monotonic()
_notify(stage_cb, "retrieving")
day = graded.get("hari_ke")
day_display = day if day is not None else '?'
if topic == "hitung_karung":
standard_block = (
"=== STANDAR MANAJEMEN PAKAN DAN SALDO KARUNG (CP 707) ===\n"
f"Umur ayam hari ke-{day_display}:\n"
"- Stok pakan: Saldo karung pakan di kandang harus selalu positif dan mencukupi kebutuhan harian ayam.\n"
"- Pencatatan: Mutasi karung masuk, saldo awal, pakan dituang, dan karung keluar harus tercatat tertib.\n"
"- Waspada: Saldo karung menipis (di bawah kebutuhan konsumsi harian) atau saldo negatif menunjukkan perlunya restock atau audit mutasi karung.\n"
"- Konsumsi air normal: 2–2.5x dari konsumsi pakan.\n"
"=== END STANDAR PAKAN ==="
)
else:
standard_block = cp707.build_cp707_standard_block(ctx) or cp707.build_book_reference_context(day if day is not None else 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_display}"
)
elif topic == "hitung_karung":
rag_query = f"manajemen stok pakan karung pakan harian konsumsi tempat pakan umur hari ke-{day_display}"
else:
rag_query = f"panduan manajemen broiler CP 707 untuk {topic} umur hari ke-{day_display}"
rag_query = f"{rag_query} {diagnose_rag_keywords(graded)}".strip()
chunks, citations = fetch_rag_chunks(rag_query, topic)
t_after_rag = time.monotonic()
_notify(stage_cb, "synthesizing")
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"
)
cited_chunks = [
f"[cp707#chunk{(citations[i].get('chunk_id') if i < len(citations) else None) or i}] {chunk}"
for i, chunk in enumerate(chunks[:4])
]
if chunks:
system_prompt += "\n[CUPLIKAN SOP CP 707 — prosa]\n" + "\n---\n".join(cited_chunks)
user_prompt = build_insight_user_prompt(
topic=topic,
kandang_name=kandang.kandang_name,
report_type=report_type,
report_period=report_period,
context=ctx,
graded=graded,
is_end_cycle=is_end_cycle,
hypotheses=hypotheses if is_end_cycle else None,
)
# Re-state output schema as the system prompt's final line (user_prompt
# repeats it again after the page-data JSON).
system_prompt += (
"\nINGAT: keluaran HANYA JSON valid sesuai FORMAT OUTPUT di atas "
"(wajib ada key kesimpulan dan insight). Tanpa teks lain."
)
# Degenerate output (schema placeholder like "..." copied verbatim) stored
# fine as JSON and rendered as an empty insight page. Reject + retry once,
# then fail loud (same policy as parse failure).
degenerate_reason: str | None = None
for attempt in range(2):
raw, usage = call_ollama(system_prompt, user_prompt)
t_after_llm = time.monotonic()
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 = enforce_bw_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:
# Fallback removed (2026-09-25): fail loudly instead of serving fake insight.
if not raw:
raise RuntimeError("Ollama tidak mengembalikan output apa pun (cek layanan LLM).")
raise RuntimeError(
"Output LLM end-cycle tidak sesuai format JSON yang diminta: "
+ repr(str(raw)[:200])
)
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:
# Fallback removed (2026-09-25): fail loudly instead of serving fake insight.
if not raw:
raise RuntimeError("Ollama tidak mengembalikan output apa pun (cek layanan LLM).")
raise RuntimeError(
"Output LLM tidak sesuai format JSON yang diminta: " + repr(str(raw)[:200])
)
parsed = enforce_mortality_wording(parsed, graded)
parsed = enforce_bw_wording(parsed, graded)
parsed = sanitize_insight_narrative(parsed)
summary = parsed["kesimpulan"]
akar_masalah = parsed.get("akar_masalah", "").strip()
clean_insight = parsed["insight"]
if akar_masalah:
daily_payload = {"akar_masalah": akar_masalah, "insight": clean_insight}
insight_text = f"{clean_insight}\n\n{DAILY_JSON_MARKER}\n{json.dumps(daily_payload, ensure_ascii=False)}"
else:
insight_text = clean_insight
check_insight = clean_insight if not is_end_cycle else insight_text
if _placeholder_text(summary) or _placeholder_text(check_insight):
degenerate_reason = f"kesimpulan={summary!r} insight={check_insight!r}"
user_prompt += (
"\nCATATAN: keluaran sebelumnya hanya placeholder. "
"Tulis kalimat lengkap sendiri untuk kesimpulan dan insight."
)
continue
degenerate_reason = None
break
if degenerate_reason:
raise RuntimeError(
"Output LLM degenerate (placeholder bukan kalimat): "
+ degenerate_reason[:200]
)
alert = alert_from_graded(graded)
_notify(stage_cb, "saving")
_log_trace(
{
"timestamp": timezone.now().isoformat(),
"request_id": str(uuid.uuid4()),
"cycle_id": cycle_id,
"kandang_name": kandang.kandang_name,
"topic": topic,
"report_type": report_type,
"report_period": report_period,
"model": getattr(settings, "LLM_MODEL_NAME", None),
"prompt_version": "v2.1",
"formulated_query": rag_query,
"retrieved_chunk_ids": [
f"cp707#{c.get('chunk_id')}" if c.get("chunk_id") is not None else f"cp707#{c.get('source')}"
for c in citations
],
"graded_alert": alert,
"llm_ok": raw is not None,
"latency_breakdown": {
"grade_ms": round((t_after_grade - t_start) * 1000),
"retrieval_ms": round((t_after_rag - t_after_grade) * 1000),
"llm_ms": round((t_after_llm - t_after_rag) * 1000),
"total_ms": round((time.monotonic() - t_start) * 1000),
},
"token_usage": usage,
}
)
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)