"""Review-screen chat: ask the run why it concluded something. Read-only by construction. The chat reads job artifacts, calls one LLM, and appends a log record; it never mutates findings, review decisions, or the report, and the prompt forbids it from emitting code or config changes. Every turn is logged twice, on purpose: - ``/review/chat_log.jsonl`` - job-local, the auditable record of what was asked about which finding and what came back. - ``REVIEW_FEEDBACK_DIR/chat_turns.jsonl`` - cross-job, append-only, so the corrections a reviewer makes in conversation ("that is not a floor drain, it is a power floor box") accumulate somewhere a future run can be primed from. Nothing reads this yet; writing it is what makes that possible later. """ import json import os import uuid from datetime import datetime, timezone from typing import Any, Dict, List, Optional from backend import config from backend.llm import call_json from backend.pipeline._serialize import dumps from backend.review.chat_context import build_context from backend.review.chat_prompts import ( REVIEW_CHAT_SYSTEM_PROMPT, REVIEW_CHAT_USER_PROMPT, ) from backend.review.feedback import append_shared_feedback _ANSWERABLE = {"yes", "partial", "no"} _ASSESSMENTS = {"looks_supported", "looks_unsupported", "cannot_tell", "not_applicable"} _CONFIDENCE = {"high", "medium", "low"} _MAX_FINDINGS = 12 _MAX_EVIDENCE = 12 class ChatError(Exception): """Raised for a caller-fixable problem (bad question, chat disabled).""" def _now() -> str: return datetime.now(timezone.utc).isoformat() def _one_of(value: Any, allowed: set, default: str) -> str: text = str(value or "").strip().lower() return text if text in allowed else default def _clean_question(raw: Any) -> str: question = str(raw or "").strip() if not question: raise ChatError("question is required") if len(question) > config.REVIEW_CHAT_MAX_QUESTION_CHARS: raise ChatError( f"question is too long (max {config.REVIEW_CHAT_MAX_QUESTION_CHARS} characters)") return question def _log_path(out_dir: str) -> str: return os.path.join(out_dir, "review", "chat_log.jsonl") def read_log(out_dir: str, review_item_id: Optional[str] = None, scope_only: bool = False) -> List[Dict[str, Any]]: """Chat turns for this job, oldest first. ``review_item_id`` filters to one finding's thread; with ``scope_only`` and no id, returns only the run-scope turns. Corrupt lines are skipped rather than failing the read - a truncated log must not hide the rest. """ path = _log_path(out_dir) if not os.path.isfile(path): return [] turns: List[Dict[str, Any]] = [] try: with open(path, encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue try: turn = json.loads(line) except json.JSONDecodeError: continue if not isinstance(turn, dict): continue if review_item_id is not None: if turn.get("review_item_id") != review_item_id: continue elif scope_only and turn.get("review_item_id") is not None: continue turns.append(turn) except OSError: return [] return turns def _append_log(out_dir: str, turn: Dict[str, Any]) -> None: """Append one turn as a JSON line; never raises on I/O failure.""" try: os.makedirs(os.path.join(out_dir, "review"), exist_ok=True) with open(_log_path(out_dir), "a", encoding="utf-8") as f: f.write(json.dumps(turn) + "\n") except OSError as e: print(f"[ReviewChat] chat log write failed: {e}") def _issue_snapshot(item: Optional[Dict]) -> Optional[Dict[str, Any]]: """The issue as it stood when asked about - the log's 'issue in question'. Copied rather than referenced by id so the log stays readable after finalization renumbers or suppresses the finding. """ if not item: return None payload = item.get("payload") or {} if item.get("kind") == "clean_cluster": return { "review_item_id": item.get("review_item_id"), "kind": item.get("kind"), "cluster_key": payload.get("key"), "location": payload.get("location"), "disciplines": payload.get("disciplines"), } return { "review_item_id": item.get("review_item_id"), "kind": item.get("kind"), "issue_id": payload.get("issue_id"), "source_stage": payload.get("source_stage"), "category": payload.get("category"), "severity": payload.get("severity"), "confidence": payload.get("confidence"), "location": payload.get("location"), "disciplines": payload.get("disciplines"), "sheets": payload.get("sheets"), "description": payload.get("description"), "blocking": item.get("blocking"), } def _history_block(turns: List[Dict[str, Any]]) -> str: if not turns: return "" recent = turns[-config.REVIEW_CHAT_HISTORY_TURNS:] lines = ["Earlier turns in this thread (oldest first):"] for turn in recent: lines.append(f"Reviewer: {turn.get('question') or ''}") lines.append(f"You: {turn.get('answer') or ''}") lines.append("") return "\n".join(lines) def _normalize_answer(parsed: Optional[Dict]) -> Optional[Dict[str, Any]]: """Coerce the model's JSON into the log/API shape, or None if unusable.""" if not isinstance(parsed, dict): return None answer = str(parsed.get("answer") or "").strip() if not answer: return None findings = [ str(item).strip() for item in (parsed.get("findings") or []) if isinstance(item, (str, int, float)) and str(item).strip() ][:_MAX_FINDINGS] evidence = [] for item in (parsed.get("evidence_cited") or [])[:_MAX_EVIDENCE]: if not isinstance(item, dict): continue evidence.append({ "artifact": str(item.get("artifact") or "").strip() or None, "sheet": item.get("sheet"), "quote": str(item.get("quote") or "").strip() or None, "why_it_matters": str(item.get("why_it_matters") or "").strip() or None, }) correction = parsed.get("suggested_category_correction") correction = str(correction).strip() if correction else "" missing = parsed.get("missing_information") return { "answer": answer, "findings": findings, "evidence_cited": evidence, "answerable": _one_of(parsed.get("answerable"), _ANSWERABLE, "partial"), "missing_information": str(missing).strip() if missing else None, "assessment_of_finding": _one_of(parsed.get("assessment_of_finding"), _ASSESSMENTS, "cannot_tell"), # The feedback signal: a reviewer correcting a misidentification in # conversation ("that is a power floor box") lands here as structured # data instead of dying in free text. "suggested_category_correction": correction or None, "confidence": _one_of(parsed.get("confidence"), _CONFIDENCE, "low"), } def _feedback_record(turn: Dict[str, Any]) -> Dict[str, Any]: """Cross-job roll-up of one turn: metadata + the correction signal. Mirrors the privacy stance of the decision labels - no images, no raw sheet dumps. The question and answer ARE carried, because a chat turn without its question is not usable as feedback; keep this store job-internal. """ issue = turn.get("issue") or {} return { "kind": "review_chat_turn", "turn_id": turn.get("turn_id"), "job_id": turn.get("job_id"), "created_at": turn.get("created_at"), "review_item_id": turn.get("review_item_id"), "scope": turn.get("scope"), "issue_id": issue.get("issue_id"), "source_stage": issue.get("source_stage"), "category": issue.get("category"), "severity": issue.get("severity"), "confidence": issue.get("confidence"), "sheets": issue.get("sheets"), "question": turn.get("question"), "answer": turn.get("answer"), "findings": turn.get("findings"), "assessment_of_finding": turn.get("assessment_of_finding"), "suggested_category_correction": turn.get("suggested_category_correction"), "answerable": turn.get("answerable"), "model": turn.get("model"), } def ask(job_id: str, out_dir: str, question: str, review_item_id: Optional[str] = None, queue: Optional[List[Dict]] = None, decisions: Optional[Dict[str, Dict]] = None) -> Dict[str, Any]: """Answer one reviewer question and log the turn. Returns the logged turn. Raises ChatError for a bad question or a disabled chat, and RuntimeError when the model call fails outright (the caller maps both to HTTP status codes). """ if not config.ENABLE_REVIEW_CHAT: raise ChatError("review chat is disabled (ENABLE_REVIEW_CHAT=false)") question = _clean_question(question) queue = queue or [] item = next((candidate for candidate in queue if candidate.get("review_item_id") == review_item_id), None) if review_item_id and item is None: raise ChatError(f"unknown review_item_id {review_item_id!r}") context = build_context(out_dir, review_item_id, queue, decisions, question) history = read_log(out_dir, review_item_id=review_item_id) if review_item_id \ else read_log(out_dir, scope_only=True) scope_line = ( f"Scope: this question is about review item {review_item_id}." if item else "Scope: this question is about the run as a whole, not one finding." ) user_text = (REVIEW_CHAT_USER_PROMPT .replace("{scope_line}", scope_line) .replace("{question}", question) .replace("{history_block}", _history_block(history)) .replace("{context}", dumps(context))) parsed = call_json( system_prompt=REVIEW_CHAT_SYSTEM_PROMPT, user_text=user_text, max_tokens=config.REVIEW_CHAT_MAX_TOKENS, model=config.REVIEW_CHAT_MODEL, usage_stage="review.chat", ) answer = _normalize_answer(parsed) if answer is None: raise RuntimeError("the model did not return a usable answer") turn = { "turn_id": uuid.uuid4().hex[:12], "job_id": job_id, "created_at": _now(), "review_item_id": review_item_id, "scope": context.get("scope"), "issue": _issue_snapshot(item), "question": question, "reviewer_decision_at_time": context.get("reviewer_decision_so_far"), "artifacts_consulted": sorted( name for name, present in (context.get("artifacts_available") or {}).items() if present ), "model": config.REVIEW_CHAT_MODEL, **answer, } _append_log(out_dir, turn) append_shared_feedback(_feedback_record(turn)) return turn def render_log_markdown(turns: List[Dict[str, Any]]) -> str: """Human-readable transcript: the issue, the questions, the findings. Grouped by review item so one finding's whole thread reads together, with run-scope questions last under their own heading. """ by_item: Dict[str, List[Dict[str, Any]]] = {} for turn in turns: by_item.setdefault(turn.get("review_item_id") or "", []).append(turn) lines = ["# Review chat log", ""] if not turns: lines.append("_No questions have been asked about this run._") return "\n".join(lines) + "\n" lines.append(f"{len(turns)} turn(s) across {len(by_item)} thread(s).") lines.append("") for item_id in sorted(by_item, key=lambda key: (key == "", key)): item_turns = by_item[item_id] issue = next((turn.get("issue") for turn in item_turns if turn.get("issue")), None) if not item_id: lines += ["## Run-scope questions", "", "_Not about a single finding._", ""] elif issue: lines.append(f"## {issue.get('issue_id') or item_id}") lines.append("") meta = [ ("Category", issue.get("category")), ("Severity", issue.get("severity")), ("Run confidence", issue.get("confidence")), ("Location", issue.get("location")), ("Sheets", ", ".join(str(s) for s in issue.get("sheets") or []) or None), ("Stage", issue.get("source_stage")), ] for label, value in meta: if value: lines.append(f"- **{label}:** {value}") if issue.get("description"): lines += ["", f"> {issue['description']}"] lines.append("") else: lines += [f"## {item_id}", ""] for turn in item_turns: lines.append(f"### Q ({turn.get('created_at') or ''})") lines += ["", turn.get("question") or "", "", "**Answer**", "", turn.get("answer") or "", ""] if turn.get("findings"): lines.append("**Findings**") lines.append("") lines += [f"- {finding}" for finding in turn["findings"]] lines.append("") if turn.get("evidence_cited"): lines += ["**Evidence cited**", ""] for item in turn["evidence_cited"]: where = item.get("artifact") or "?" sheet = f" ({item['sheet']})" if item.get("sheet") else "" quote = item.get("quote") or "" lines.append(f"- `{where}`{sheet}: \"{quote}\"") if item.get("why_it_matters"): lines.append(f" - {item['why_it_matters']}") lines.append("") tail = [ ("Answerable", turn.get("answerable")), ("Assessment", turn.get("assessment_of_finding")), ("Confidence", turn.get("confidence")), ("Missing", turn.get("missing_information")), ("Suggested correction", turn.get("suggested_category_correction")), ("Model", turn.get("model")), ] lines.append(" | ".join(f"{label}: {value}" for label, value in tail if value)) lines.append("") return "\n".join(lines) + "\n"