526 lines
18 KiB
Python
526 lines
18 KiB
Python
"""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
|
||
|
||
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 MORTALITAS (WAJIB DIPATUHI — TIDAK BOLEH DILANGGAR):
|
||
- Standar CP 707: mortalitas kumulatif NORMAL adalah < 5%.
|
||
- Mortalitas >= 5% dan <= 7% = TINGGI (warning) — WAJIB disebut "TINGGI", bukan "rendah" atau "normal".
|
||
- Mortalitas > 7% = SANGAT TINGGI (critical) — WAJIB disebut "SANGAT TINGGI".
|
||
- Mortalitas < 5% = rendah/normal.
|
||
- DILARANG KERAS menyebut mortalitas sebagai "rendah" atau "normal" jika nilainya >= 5%.
|
||
|
||
ATURAN ANGKA & ARAH (WAJIB DIPATUHI):
|
||
- Setiap angka yang Anda tulis HARUS ada di [Data Halaman (JSON)] atau di blok STANDAR CP 707
|
||
atau di [GRADED FACTS]. DILARANG menghitung sendiri, memperkirakan, atau membalik arah tren.
|
||
- 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.
|
||
|
||
FORMAT OUTPUT (JSON SAJA):
|
||
{"kesimpulan":"...","insight":"..."}
|
||
"""
|
||
|
||
|
||
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)
|
||
if report_type == "end_cycle":
|
||
# Deterministic weekly KPI rollup when weekly_summaries absent.
|
||
if "weekly_summaries" not in out:
|
||
out["weekly_summaries"] = _build_weekly_summaries(out)
|
||
out["catatan_end_cycle"] = (
|
||
"Ringkasan KPI per minggu dihitung di backend. Narasikan dari weekly_summaries + graded_facts; "
|
||
"jangan menghitung ulang."
|
||
)
|
||
out["report_type"] = report_type
|
||
out["report_period"] = report_period
|
||
return out
|
||
|
||
|
||
def _build_weekly_summaries(context: dict[str, Any]) -> list[dict[str, Any]]:
|
||
histories = []
|
||
for key in ("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 = row.get("hari") or row.get("day")
|
||
try:
|
||
day_i = int(day)
|
||
except (TypeError, ValueError):
|
||
continue
|
||
week = max(1, (day_i + 6) // 7)
|
||
by_week.setdefault(week, []).append(row)
|
||
|
||
summaries = []
|
||
for week, rows in sorted(by_week.items()):
|
||
fcrs = [r.get("fcr_aktual") or r.get("fcr") for r in rows if isinstance(r.get("fcr_aktual") or r.get("fcr"), (int, float))]
|
||
summaries.append(
|
||
{
|
||
"minggu_ke": week,
|
||
"jumlah_hari_data": len(rows),
|
||
"fcr_rata_rasio": round(sum(fcrs) / len(fcrs), 3) if fcrs else None,
|
||
"fcr_akhir_rasio": fcrs[-1] if fcrs else None,
|
||
}
|
||
)
|
||
return summaries
|
||
|
||
|
||
def grade_context(context: dict[str, Any]) -> dict[str, Any]:
|
||
analysis = cp707.analyze_with_cp707_standards(context)
|
||
graded: dict[str, Any] = {
|
||
"kandangId": context.get("kandangId") or context.get("kandang_id"),
|
||
"hari_ke": analysis.get("dayAge") or _day_age(context),
|
||
"analyses": analysis.get("analyses") or {},
|
||
}
|
||
# 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 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 ""
|
||
)
|
||
insight = (
|
||
parsed.get("insight")
|
||
or parsed.get("rekomendasi")
|
||
or parsed.get("insights")
|
||
or ""
|
||
)
|
||
if not kesimpulan and not insight:
|
||
return None
|
||
return {
|
||
"kesimpulan": str(kesimpulan) or "Model tidak mengembalikan kesimpulan eksplisit.",
|
||
"insight": str(insight) or "Model tidak mengembalikan rekomendasi eksplisit.",
|
||
}
|
||
|
||
|
||
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],
|
||
}
|
||
|
||
|
||
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, source_override=AIInsight.SOURCE_CACHE)
|
||
|
||
|
||
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, source_override: str | None = None) -> dict[str, Any]:
|
||
norm_citations = decode_citations(row.citations)
|
||
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,
|
||
"source": source_override or row.source or AIInsight.SOURCE_GENERATED,
|
||
"summary": row.summary or "",
|
||
"insight": row.insight_text or "",
|
||
"insight_text": "\n\n".join(
|
||
part for part in [row.summary or "", row.insight_text or ""] if part
|
||
)
|
||
or row.insight_text
|
||
or "",
|
||
"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)
|
||
|
||
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"
|
||
f"{standard_block}\n"
|
||
f"[GRADED FACTS]\n{json.dumps(graded, ensure_ascii=False, default=str)}\n"
|
||
)
|
||
if chunks:
|
||
system_prompt += "\n[CUPLIKAN SOP CP 707 — prosa]\n" + "\n---\n".join(chunks[:4])
|
||
|
||
user_prompt = (
|
||
f"Buat insight topik `{topic}` untuk kandang `{kandang.kandang_name}` "
|
||
f"(id={kandang_id}), periode `{report_type}/{report_period}`.\n"
|
||
f"[Data Halaman (JSON)]\n{json.dumps(ctx, ensure_ascii=False, default=str)}"
|
||
)
|
||
|
||
raw = call_ollama(system_prompt, user_prompt)
|
||
parsed = parse_llm_json(raw or "")
|
||
source = AIInsight.SOURCE_GENERATED
|
||
if not parsed:
|
||
parsed = local_fallback_insight(graded, topic)
|
||
source = AIInsight.SOURCE_LOCAL_FALLBACK
|
||
|
||
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,
|
||
"source": source,
|
||
"citations": encode_citations(citations),
|
||
},
|
||
)
|
||
return serialize_insight(row, source_override=source)
|