"""OpenAI tool-calling loop for the RICS inspector: the model chooses tools and order. Unlike the fixed gather→draft pipeline, this module runs a multi-turn ``chat.completions`` session with ``tools=…`` until the model completes the **server-enforced** workflow (extraction audit → section plan → optional Section C rating summary when the active product pack includes condition ratings → ``submit_inspection_section``), or a round limit is reached. """ from __future__ import annotations import asyncio import json import logging import re from typing import Any from openai import OpenAI from app.agentic import tools as agent_tools from app.config import settings from app.generator.postprocess import ( _L1_PLACEHOLDER, async_enforce_verify, enforce_verify, strip_l1_advice, strip_l1_advice_payload, ) from app.models.schemas import SearchResult, WritingStyleProfile from app.services.provenance_enrichment import fetch_doc_filenames from app.services.generation import _ai_level_to_params from app.templates.registry import SurveyTemplatePack, get_survey_pack, get_template from .models import EvidenceItem, RiskItem, StructuredReport logger = logging.getLogger(__name__) _inspector_async_client: Any = None def _inspector_async_client() -> Any: """Reuse one AsyncOpenAI client across inspector tool rounds.""" global _inspector_async_client if _inspector_async_client is None: from openai import AsyncOpenAI _inspector_async_client = AsyncOpenAI(api_key=settings.openai_api_key) return _inspector_async_client def _strip_missing_fact_phrase(text: str) -> str: """Final defensive cleanup for legacy missing-data filler phrases.""" if not isinstance(text, str) or not text: return text if isinstance(text, str) else "" cleaned = re.sub( r"\bInformation not provided in source document\b[.,;:!?]*", "", text, flags=re.IGNORECASE, ) cleaned = re.sub( r"\bWe were unable to verify this during inspection\b[.,;:!?]*", "", cleaned, flags=re.IGNORECASE, ) cleaned = re.sub(r"\s{2,}", " ", cleaned) cleaned = re.sub(r"\s+([,.;:])", r"\1", cleaned) return cleaned.strip() async def _enforce_verify_submit_payload( payload: dict[str, str], *, bullets: list[str], snippets: list[str], peer_sections: dict[str, str] | None, pinned_identity: dict[str, str] | None = None, ) -> dict[str, str]: """Run the non-invention guard on the five free-text fields the LLM submitted. The agentic ``submit_inspection_section`` tool accepts five strings drafted by the LLM. The legacy fixed pipeline runs ``async_enforce_verify`` on every drafted section before persistence; the agentic path bypassed it entirely, which is how invented addresses, postcodes, and construction-type claims were reaching the user-visible report. We treat all of the following as "trusted source material" — anything the LLM uses outside of these is invention: - ``bullets`` : the surveyor's messy notes for THIS section - ``peer_sections`` : the surveyor's draft notes for sibling sections (legitimate cross-section context) - ``snippets`` : retrieved tenant + KB chunks deduped by the loop The five fields are verified concurrently with ``asyncio.gather``; each call runs the cheap regex pass first and only falls through to the LLM grounding pass when an OpenAI key is configured. Cost is at most one extra mini-LLM call per non-empty field (~£0.0001–0.0005 each). """ peer_values = list((peer_sections or {}).values()) full_bullets = list(bullets) + [v for v in peer_values if isinstance(v, str) and v.strip()] field_names = ( "executive_summary", "property_description", "condition_assessment", "defects_and_risks", "recommendations", ) async def _verify_one(name: str) -> tuple[str, str]: text = payload.get(name) or "" if not text.strip(): return name, text try: verified = await async_enforce_verify( text=text, bullets=full_bullets, snippets=snippets, pinned_identity=pinned_identity, openai_api_key=settings.openai_api_key or "", model=settings.chat_model, ) return name, _strip_missing_fact_phrase(verified) except Exception as exc: # noqa: BLE001 # Grounding must never block the section — a verifier failure is # logged and the original LLM text is returned, matching the # behaviour of the legacy pipeline. logger.warning( "Non-invention guard failed for field=%s (%s); returning raw text", name, exc, ) return name, _strip_missing_fact_phrase(text) results = await asyncio.gather(*(_verify_one(n) for n in field_names)) return {**payload, **dict(results)} # --------------------------------------------------------------------------- # Comprehensive non-invention guard for ALL LLM-emitted artifacts. # --------------------------------------------------------------------------- # The agentic inspector emits text in *six* places, not five: # # 1. submit_inspection_section payload (5 user-visible body fields) — full # regex + LLM grounding pass via _enforce_verify_submit_payload. # 2. submit_extraction_audit (property_address, surveyor_name, …) # 3. submit_section_plan (claims[*].claim, non_claims[*]) # 4. submit_condition_rating_summary (justification, summary_markdown_table) # 5. RiskAssessmentAgent.assess (risk[*].risk, risk[*].action) # 6. compliance check (heuristic — already deterministic) # # Without (2)–(5) the user could see invented postcodes, surveyor names, or # risk-action references in the metadata bundle the API returns # (`/agentic/...` endpoints surface extraction_audit + section_plan + # condition_rating_summary in inspector_meta). The five body fields run the # expensive 2-layer guard because they ARE the report; auxiliary fields run # the cheap synchronous regex pass — that catches postcodes, named entities, # and bare numbers without adding 30+ extra LLM calls per section. def _regex_verify_str( text: Any, *, bullets: list[str], snippets: list[str], pinned_identity: dict[str, str] | None = None, ) -> str: """Run the synchronous regex-only guard on a single string-typed field. Returns the input unchanged when it isn't a non-empty string. Used for auxiliary fields where latency / cost matters and most violations are short factual claims (postcodes, addresses, named persons/firms). """ if not isinstance(text, str) or not text.strip(): return text if isinstance(text, str) else "" return enforce_verify( text=text, bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ) def _verify_audit_payload( audit: Any, *, bullets: list[str], snippets: list[str], pinned_identity: dict[str, str] | None = None, ) -> dict[str, Any] | None: if not isinstance(audit, dict): return audit out = dict(audit) for key in ( "property_address", "surveyor_name", "rics_number", "company_name", "inspection_date", "report_reference", "property_type", "construction_details", "services", ): if key in out: out[key] = _regex_verify_str( out.get(key), bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ) for list_key in ("client_names", "observed_defects"): items = out.get(list_key) if isinstance(items, list): out[list_key] = [ _regex_verify_str( x, bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ) for x in items if isinstance(x, str) ] return out def _verify_plan_payload( plan: Any, *, bullets: list[str], snippets: list[str], pinned_identity: dict[str, str] | None = None, ) -> dict[str, Any] | None: if not isinstance(plan, dict): return plan out = dict(plan) claims = out.get("claims") if isinstance(claims, list): verified_claims: list[dict[str, Any]] = [] for c in claims: if isinstance(c, dict): cc = dict(c) if "claim" in cc: cc["claim"] = _regex_verify_str( cc.get("claim"), bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ) verified_claims.append(cc) else: verified_claims.append(c) out["claims"] = verified_claims nc = out.get("non_claims") if isinstance(nc, list): out["non_claims"] = [ _regex_verify_str( x, bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ) for x in nc if isinstance(x, str) ] return out def _verify_condition_payload( cond: Any, *, bullets: list[str], snippets: list[str], pinned_identity: dict[str, str] | None = None, ) -> dict[str, Any] | None: if not isinstance(cond, dict): return cond out = dict(cond) for k in ("justification", "summary_markdown_table"): if k in out: out[k] = _regex_verify_str( out.get(k), bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ) return out def _verify_risks( risks: list[RiskItem], *, bullets: list[str], snippets: list[str], pinned_identity: dict[str, str] | None = None, ) -> list[RiskItem]: """Strip invented postcodes / addresses / named firms from LLM-drafted risks. `RiskItem` is frozen, so we rebuild each item; `category`, `severity`, and `likelihood` are constrained to known enum values upstream and never need verification, so we only sweep `risk` and `action`. """ out: list[RiskItem] = [] for r in risks: try: out.append( RiskItem( category=r.category, risk=_regex_verify_str( r.risk, bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ), severity=r.severity, likelihood=r.likelihood, action=_regex_verify_str( r.action, bullets=bullets, snippets=snippets, pinned_identity=pinned_identity, ), evidence=r.evidence, ) ) except Exception as exc: # noqa: BLE001 logger.warning("Risk verification failed (%s); keeping raw item", exc) out.append(r) return out def _openai_inspector_tool_definitions() -> list[dict[str, Any]]: """Narrow tool schemas: tenant/document IDs are injected server-side.""" return [ { "type": "function", "function": { "name": "retrieve_survey_rag", "description": ( "Search the tenant's indexed RAG corpus (uploaded survey PDF/DOCX and reference " "library), prioritising the report's primary document. Call this when you need " "verbatim or technical evidence from the uploaded survey before writing findings." ), "parameters": { "type": "object", "properties": { "query": { "type": "string", "description": "Semantic query distilled from the inspection notes (concise).", }, "k": { "type": "integer", "description": "Max chunks to retrieve before reranking (default 14).", "minimum": 1, "maximum": 30, }, "rerank_top_n": { "type": "integer", "description": "How many chunks to return after reranking (default 7).", "minimum": 1, "maximum": 15, }, }, "required": ["query"], }, }, }, { "type": "function", "function": { "name": "retrieve_rics_kb", "description": ( "Search the local RICS standards / exemplar knowledge base for professional " "wording, definitions, and compliance framing. Use when notes are thin or you " "need authoritative context (still ground factual claims in survey RAG)." ), "parameters": { "type": "object", "properties": { "query": {"type": "string"}, "k": {"type": "integer", "description": "Max KB chunks to fetch (default 10).", "minimum": 1, "maximum": 20}, "rerank_top_n": {"type": "integer", "description": "Chunks to keep after rerank (default 5).", "minimum": 1, "maximum": 10}, "hierarchy_level": { "type": "string", "enum": ["document", "section", "paragraph"], "description": "Optional granularity filter when the index supports hierarchy.", }, }, "required": ["query"], }, }, }, { "type": "function", "function": { "name": "scan_note_duplicates", "description": ( "Check the tenant's indexed library for semantically similar passages and " "detect near-duplicate wording across other sections' drafts. Use when notes " "may overlap prior reports or other sections." ), "parameters": { "type": "object", "properties": { "text": {"type": "string", "description": "Notes block to compare (e.g. joined bullets)."}, "section_code": {"type": "string"}, "peer_sections": { "type": "object", "additionalProperties": {"type": "string"}, "description": "Map section_code → draft text for cross-section overlap.", }, "exclude_document_ids": { "type": "array", "items": {"type": "string"}, "description": "Upload UUIDs to omit from library similarity (e.g. primary survey).", }, }, "required": ["text"], }, }, }, { "type": "function", "function": { "name": "submit_extraction_audit", "description": ( "MANDATORY (server-enforced before final submit): structured extraction from the " "inspection notes and survey RAG snippets — what is evidenced vs unknown. " "Call once per section after you have run retrieve_survey_rag / retrieve_rics_kb " "enough to understand the fact basis." ), "parameters": { "type": "object", "properties": { "property_address": {"type": "string"}, "client_names": {"type": "array", "items": {"type": "string"}}, "surveyor_name": {"type": "string"}, "rics_number": {"type": "string"}, "company_name": {"type": "string"}, "inspection_date": {"type": "string"}, "report_reference": {"type": "string"}, "property_type": {"type": "string"}, "construction_details": {"type": "string"}, "services": {"type": "string"}, "observed_defects": {"type": "array", "items": {"type": "string"}}, "limitations_to_inspection": {"type": "array", "items": {"type": "string"}}, "assumptions": {"type": "array", "items": {"type": "string"}}, }, "required": [ "property_address", "client_names", "surveyor_name", "rics_number", "company_name", "inspection_date", "report_reference", "property_type", "construction_details", "services", "observed_defects", "limitations_to_inspection", "assumptions", ], }, }, }, { "type": "function", "function": { "name": "submit_section_plan", "description": ( "MANDATORY (server-enforced before final submit): a short plan mapping claims to " "evidence chunk_ids / doc_ids you intend to rely on. Ratings (1/2/3/NI) must appear " "only where the active RICS product tier uses condition ratings for that element group." ), "parameters": { "type": "object", "properties": { "claims": { "type": "array", "items": { "type": "object", "properties": { "claim": {"type": "string"}, "evidence_refs": { "type": "array", "items": {"type": "string"}, "description": "chunk_id values and/or doc_id:chunk_id pairs from tool hits.", }, "risk": { "type": "string", "enum": ["low", "medium", "high"], }, }, "required": ["claim", "evidence_refs", "risk"], }, }, "non_claims": { "type": "array", "items": {"type": "string"}, "description": "Boilerplate you will keep generic (no site-specific facts).", }, }, "required": ["claims", "non_claims"], }, }, }, { "type": "function", "function": { "name": "submit_condition_rating_summary", "description": ( "When the active survey product includes condition ratings anywhere in its template " "pack, Section C must summarise ratings in a markdown table before the final submit. " "If ratings truly do not apply for this property type, set ratings_not_applicable=true " "with a one-line justification instead of inventing counts." ), "parameters": { "type": "object", "properties": { "ratings_not_applicable": {"type": "boolean"}, "justification": {"type": "string"}, "summary_markdown_table": { "type": "string", "description": "Markdown table with columns Category | Count NI | Count 3 | Count 2 | Count 1", }, }, "required": ["ratings_not_applicable"], }, }, }, { "type": "function", "function": { "name": "submit_inspection_section", "description": ( "Submit the final drafted content for this RICS section. Server will reject this " "until submit_extraction_audit and submit_section_plan have been accepted. " "When the product pack includes condition ratings, Section C also requires " "submit_condition_rating_summary first. " "Ground factual statements in retrieve_survey_rag / retrieve_rics_kb results. " "If information is missing, omit the claim rather than writing placeholder " "sentences. Never insert missing-data text mid-sentence. " "UK English, MRICS tone. Call alone (no other tools in the same assistant message). " "PRODUCE PROPERTY-SPECIFIC PROSE AT THE TIER-APPROPRIATE LENGTH — NOT SUMMARIES, " "NOR PADDING. Each free-text field below is the user-visible report; aim at the " "per-tier word_targets in the SYSTEM PROMPT (L1 condensed observation, L2 " "proportionate buyer advice, L3 diagnostic depth). Do NOT collapse an L3 field to " "one or two sentences; do NOT inflate an L1 field beyond its observation-only " "range. Cite specific defects, locations, materials, and observed evidence — " "never generic boilerplate." ), "parameters": { "type": "object", # Schema descriptions deliberately defer length / depth to # the per-tier SYSTEM PROMPT (word_targets + depth_directive # in :func:`_inspector_system_prompt`) instead of baking in # L3-leaning ranges. Otherwise the model sees "target # 200–500 words" for condition_assessment and tries to hit # that even when the active tier is L1 (40–90 words), or # tries a Level-3-style cause/implications/options layering # in an L1 product (which is forbidden). "properties": { "executive_summary": { "type": "string", "description": ( "Headline summary for THIS section. Length per tier — see SYSTEM " "PROMPT word_targets. State scope of inspection, principal " "observations, and overall outcome. Avoid generic openings like " "'This section covers…'." ), }, "property_description": { "type": "string", "description": ( "Specific factual description of the section subject. Length per " "tier — see SYSTEM PROMPT word_targets. Cover construction type, " "materials, accommodation, age, services, and dimensions WHEN " "supported by bullets/snippets. Use the approved missing-info " "phrase only for facts truly unavailable." ), }, "condition_assessment": { "type": "string", "description": ( "Condition narrative. Length and depth per tier — see SYSTEM " "PROMPT word_targets and depth_directive (L1 observation-only; " "L2 proportionate; L3 cause → implications → options). Cite each " "defect with location and severity at the depth appropriate to " "the active tier; do not over-extend an L1 paragraph or under-fill " "an L3 one." ), }, "defects_and_risks": { "type": "string", "description": ( "Defects mapped to category (structural / moisture / electrical / " "etc.), with location, cause (where supported), and consequence. " "Length per tier — see SYSTEM PROMPT word_targets. List each " "evidenced defect; do not omit any from the bullets." ), }, "recommendations": { "type": "string", "description": ( "Actionable recommendations with urgency, tied to an observed " "defect or evidence reference. Length per tier — see SYSTEM " "PROMPT word_targets. For Level 1 (Condition Report) the system " "prompt forbids advice — in that case write only that no " "recommendations are issued at this tier, and stop." ), }, }, "required": [ "executive_summary", "property_description", "condition_assessment", "defects_and_risks", "recommendations", ], }, }, }, ] def _compact_hits_for_llm(hits: list[SearchResult], *, max_items: int = 8, max_chars: int = 520) -> list[dict[str, Any]]: out: list[dict[str, Any]] = [] for r in hits[:max_items]: t = (r.text or "").strip().replace("\n", " ") if len(t) > max_chars: t = t[: max_chars - 1] + "…" out.append( { "chunk_id": r.chunk_id, "doc_id": r.doc_id, "score": round(float(r.score), 5), "kb": bool(getattr(r, "kb", False)), "text": t, } ) return out def _pack_has_condition_ratings(pack: SurveyTemplatePack) -> bool: return any(bool(t.has_condition_rating) for t in pack._by_code.values()) def _section_c_requires_rating_summary(*, section_code: str, pack: SurveyTemplatePack) -> bool: if section_code != "C": return False return _pack_has_condition_ratings(pack) def _norm_str_list(val: Any, *, min_items: int, max_items: int, max_item_len: int) -> tuple[str, list[str] | None]: if not isinstance(val, list): return ("must be a JSON array of strings", None) out: list[str] = [] for x in val[:max_items]: s = str(x).strip() if not s: continue if len(s) > max_item_len: s = s[: max_item_len - 1] + "…" out.append(s) if len(out) < min_items: return (f"need at least {min_items} non-empty string item(s)", None) return ("", out) def _validate_extraction_audit(args: dict[str, Any]) -> tuple[str, dict[str, Any] | None]: # Mandatory extraction schema for reconstruction. Use the two approved phrases when unknown. required_str_fields = [ "property_address", "surveyor_name", "rics_number", "company_name", "inspection_date", "report_reference", "property_type", "construction_details", "services", ] out: dict[str, Any] = {} errs: list[str] = [] for k in required_str_fields: v = str(args.get(k, "")).strip() if not v: errs.append(f"{k} is required") elif len(v) > 1200: out[k] = v[:1199] + "…" else: out[k] = v cn = args.get("client_names") if not isinstance(cn, list): errs.append("client_names must be a JSON array of strings") client_names = None else: client_names = [str(x).strip() for x in cn if str(x).strip()][:12] if not client_names: errs.append("client_names must contain at least one name or an approved missing-info phrase") if client_names is not None: out["client_names"] = client_names od_msg, observed_defects = _norm_str_list(args.get("observed_defects"), min_items=0, max_items=40, max_item_len=420) if od_msg: errs.append(f"observed_defects: {od_msg}") else: out["observed_defects"] = observed_defects lim_msg, limitations = _norm_str_list(args.get("limitations_to_inspection"), min_items=0, max_items=24, max_item_len=420) if lim_msg: errs.append(f"limitations_to_inspection: {lim_msg}") else: out["limitations_to_inspection"] = limitations asm_msg, assumptions = _norm_str_list(args.get("assumptions"), min_items=0, max_items=16, max_item_len=420) if asm_msg: errs.append(f"assumptions: {asm_msg}") else: out["assumptions"] = assumptions if errs: return ("; ".join(errs), None) return ("", out) def _validate_section_plan(args: dict[str, Any]) -> tuple[str, dict[str, Any] | None]: claims_raw = args.get("claims") if not isinstance(claims_raw, list) or not claims_raw: return ("claims must be a non-empty JSON array", None) claims_out: list[dict[str, Any]] = [] for c in claims_raw[:30]: if not isinstance(c, dict): return ("each claim must be an object", None) claim = str(c.get("claim", "")).strip() if len(claim) < 12: return ("each claim.claim must be a substantive string", None) refs = c.get("evidence_refs") if not isinstance(refs, list) or not refs: return ("each claim needs a non-empty evidence_refs array", None) refs_s = [str(x).strip() for x in refs if str(x).strip()][:24] if not refs_s: return ("evidence_refs must contain non-empty strings", None) risk = str(c.get("risk", "")).strip().lower() if risk not in ("low", "medium", "high"): return ("claim.risk must be one of: low | medium | high", None) claims_out.append({"claim": claim[:900], "evidence_refs": refs_s, "risk": risk}) nc_msg, non_claims = _norm_str_list(args.get("non_claims"), min_items=0, max_items=24, max_item_len=300) if nc_msg: return (nc_msg, None) assert non_claims is not None return ("", {"claims": claims_out, "non_claims": non_claims}) def _validate_condition_rating_summary(args: dict[str, Any]) -> tuple[str, dict[str, Any] | None]: na = args.get("ratings_not_applicable") if not isinstance(na, bool): return ("ratings_not_applicable must be a boolean", None) justification = str(args.get("justification", "")).strip() table = str(args.get("summary_markdown_table", "")).strip() if na: if len(justification) < 12: return ("when ratings_not_applicable=true, justification must explain why (substantive).", None) return ("", {"ratings_not_applicable": True, "justification": justification[:800], "summary_markdown_table": ""}) lowered = table.lower() if "|" not in table or "\n" not in table: return ("summary_markdown_table must be a markdown table (header + separator + rows).", None) for token in ("ni", "3", "2", "1"): if token not in lowered: return ("table must include NI and numeric rating columns/labels for 3, 2, and 1.", None) return ("", {"ratings_not_applicable": False, "justification": "", "summary_markdown_table": table[:6000]}) async def _dispatch_tool( *, db, tenant_id: str, primary_document_id: str, reference_document_ids: list[str] | None, default_hierarchy_level: str | None, peer_sections_default: dict[str, str] | None, section_code: str, pack: SurveyTemplatePack, gate: dict[str, Any], name: str, args: dict[str, Any], ) -> tuple[str, list[SearchResult]]: """Run one tool; return JSON string for the assistant plus any new search hits.""" extra_hits: list[SearchResult] = [] try: if name == "retrieve_survey_rag": q = str(args.get("query", "")).strip() if not q: return json.dumps({"error": "empty query"}), [] k = int(args.get("k", 14)) rn = int(args.get("rerank_top_n", 7)) k = max(4, min(48, k)) rn = max(2, min(16, rn)) hits = await agent_tools.retrieve_tenant_evidence_async( query=q, tenant_id=tenant_id, primary_document_id=primary_document_id, secondary_document_ids=reference_document_ids, k=k, rerank_top_n=rn, ) extra_hits.extend(hits) return json.dumps({"hits": _compact_hits_for_llm(hits)}), hits if name == "retrieve_rics_kb": from app.services.personalised_rag import tenant_has_personal_library if settings.personalised_style_rag_enabled and await tenant_has_personal_library( tenant_id ): return ( json.dumps( { "hits": [], "message": ( "KB retrieval disabled: use retrieve_survey_rag for this " "tenant's private uploaded report library." ), } ), [], ) q = str(args.get("query", "")).strip() if not q: return json.dumps({"hits": [], "message": "KB disabled or empty query"}), [] k = int(args.get("k", 10)) rn = int(args.get("rerank_top_n", 5)) hl = args.get("hierarchy_level") if isinstance(hl, str) and hl not in ("document", "section", "paragraph"): hl = default_hierarchy_level elif hl is None: hl = default_hierarchy_level if default_hierarchy_level in ("document", "section", "paragraph") else None hits = await agent_tools.retrieve_kb_guidance_async( query=q, k=k, hierarchy_level=hl, rerank_top_n=rn ) extra_hits.extend(hits) return json.dumps({"hits": _compact_hits_for_llm(hits), "kb_enabled": settings.knowledge_base_enabled}), hits if name == "scan_note_duplicates": text = str(args.get("text", "")).strip() if not text: return json.dumps({"library_matches": [], "draft_overlaps": [], "message": "no text"}), [] peers = args.get("peer_sections") if not isinstance(peers, dict): peers = dict(peer_sections_default or {}) sec = args.get("section_code") sec_s = str(sec) if sec is not None else None excl = args.get("exclude_document_ids") excl_l = [str(x) for x in excl] if isinstance(excl, list) else [] sim = await agent_tools.find_similar_library_and_peers( db, tenant_id, text=text, section_code=sec_s, peer_sections=peers, exclude_document_ids=excl_l, ) payload = { "library_matches": [m.model_dump() for m in sim.library_matches[:12]], "draft_overlaps": [m.model_dump() for m in sim.draft_overlaps[:12]], "message": (sim.message or "")[:500], } return json.dumps(payload), [] if name == "submit_extraction_audit": err, payload = _validate_extraction_audit(args) if err or payload is None: return json.dumps({"ok": False, "error": err or "invalid payload"}), [] gate["audit_payload"] = payload return json.dumps({"ok": True, "accepted": "extraction_audit"}), [] if name == "submit_section_plan": err, payload = _validate_section_plan(args) if err or payload is None: return json.dumps({"ok": False, "error": err or "invalid payload"}), [] if gate.get("audit_payload") is None: return json.dumps( { "ok": False, "error": "submit_extraction_audit must be accepted before submit_section_plan.", } ), [] gate["plan_payload"] = payload return json.dumps({"ok": True, "accepted": "section_plan"}), [] if name == "submit_condition_rating_summary": err, payload = _validate_condition_rating_summary(args) if err or payload is None: return json.dumps({"ok": False, "error": err or "invalid payload"}), [] if not _section_c_requires_rating_summary(section_code=section_code, pack=pack): return json.dumps( { "ok": False, "error": "submit_condition_rating_summary is only used for Section C when the active product pack includes condition ratings.", } ), [] # Enforce RICS-like grouping: the table must mention every rated element group prefix # present in this product pack (e.g. E/F/G for Building Survey). required_prefixes = sorted( {t.code[0] for t in pack._by_code.values() if t.has_condition_rating and t.code and t.code[0].isalpha()} ) table = str(payload.get("summary_markdown_table") or "") if required_prefixes and table: missing_prefixes = [] for pfx in required_prefixes: # row like "| E1" or "| E " if f"| {pfx}" not in table and f"|{pfx}" not in table: missing_prefixes.append(pfx) if missing_prefixes: return json.dumps( { "ok": False, "error": ( "Condition rating summary table must include rows for these element groups " f"(from the active template pack): {', '.join(missing_prefixes)}." ), } ), [] gate["condition_payload"] = payload return json.dumps({"ok": True, "accepted": "condition_rating_summary"}), [] if name == "submit_inspection_section": if gate.get("audit_payload") is None: return json.dumps( {"ok": False, "error": "Call and pass submit_extraction_audit before submit_inspection_section."} ), [] if gate.get("plan_payload") is None: return json.dumps( {"ok": False, "error": "Call and pass submit_section_plan before submit_inspection_section."} ), [] if _section_c_requires_rating_summary(section_code=section_code, pack=pack) and gate.get("condition_payload") is None: return json.dumps( { "ok": False, "error": "Section C in this product tier requires submit_condition_rating_summary before submit_inspection_section.", } ), [] gate["submit_payload"] = { "executive_summary": str(args.get("executive_summary", "")).strip(), "property_description": str(args.get("property_description", "")).strip(), "condition_assessment": str(args.get("condition_assessment", "")).strip(), "defects_and_risks": str(args.get("defects_and_risks", "")).strip(), "recommendations": str(args.get("recommendations", "")).strip(), } return json.dumps({"ok": True, "accepted": "submit_inspection_section"}), [] except Exception as exc: # noqa: BLE001 logger.exception("inspector tool %s failed", name) return json.dumps({"error": str(exc)}), [] return json.dumps({"error": f"unknown tool {name}"}), [] async def _evidence_items_async(db, tenant_id: str, uniq: list[SearchResult]) -> list[EvidenceItem]: tenant_doc_ids = {r.doc_id for r in uniq if r.tenant_id == tenant_id and r.doc_id} filenames: dict[str, str | None] = {} if db is not None and tenant_doc_ids: filenames = await fetch_doc_filenames(db, tenant_id, tenant_doc_ids) items: list[EvidenceItem] = [] for r in uniq: is_kb = bool(getattr(r, "kb", False)) or r.tenant_id == settings.knowledge_base_tenant_id src = (getattr(r, "source", None) or None) if is_kb else filenames.get(r.doc_id) if not src and is_kb: src = "Local RICS knowledge base" items.append( EvidenceItem( doc_id=r.doc_id, chunk_id=r.chunk_id, score=float(r.score), text=r.text, source=str(src) if src else None, section_hint=getattr(r, "section_title", None), kb=is_kb, ) ) return items def _inspector_system_prompt( *, tenant_id: str, primary_document_id: str, survey_level: int | None, pack: SurveyTemplatePack, section_code: str, section_title: str, style_profile: WritingStyleProfile | None, identity_block: str | None = None, ) -> str: tone = style_profile.tone if style_profile and style_profile.tone else "professional, precise UK building surveyor" order_preview = ", ".join(list(pack.section_order)[:34]) if len(pack.section_order) > 34: order_preview += ", …" ratings_pack = _pack_has_condition_ratings(pack) c_extra = "" if section_code == "C" and ratings_pack: c_extra = ( " For Section C in this product tier, you must also call submit_condition_rating_summary " "(markdown rating table, or ratings_not_applicable with justification) before the final submit." ) # ── Per-tier depth directive ─────────────────────────────────────────── # Real RICS L3 element narratives describe each defect as # cause → implications → options/next steps, in continuous diagnostic # prose. Telling the model the *upper* end of the word range (and that # L3 specifically must layer cause/implications/options) is what moves # gpt-4o-mini from "150-word HomeBuyer-style paragraph" to actual # diagnostic depth. Without an explicit L3 directive the model defaults # to its trained-average paragraph length, which is roughly L2. try: lvl = int(survey_level if survey_level is not None else pack.level) except Exception: # noqa: BLE001 lvl = pack.level if lvl >= 3: depth_directive = ( "DEPTH FOR LEVEL 3 (BUILDING SURVEY — DIAGNOSTIC MODE): aim at the UPPER end of every " "word range below. For each material defect, weave (within the relevant field): (a) what " "was observed, (b) the likely cause/mechanism (only if supported by evidence), " "(c) implications/risks if unaddressed, and (d) options/next steps. Do NOT collapse " "into HomeBuyer-style brevity. Each defect must surface in defects_and_risks AND " "carry through to recommendations with a tied action. " ) word_targets = ( "executive_summary ~30–60 words, property_description ~35–70 words, " "condition_assessment ~60–120 words, defects_and_risks ~35–80 words, " "recommendations ~20–45 words. " ) elif lvl == 2: depth_directive = ( "DEPTH FOR LEVEL 2 (HomeBuyer): proportionate buyer-focused advice. Cover practical next " "steps for material defects (e.g. obtain quotations, further checks). Avoid deep " "diagnostic speculation unless the bullets/snippets explicitly support it. " ) word_targets = ( "executive_summary ~20–45 words, property_description ~25–50 words, " "condition_assessment ~45–90 words, defects_and_risks ~25–55 words, " "recommendations ~15–35 words. " ) else: depth_directive = ( "DEPTH FOR LEVEL 1 (Condition Report — observation mode): record condition concisely. " "Do NOT advise on repairs, recommend works, or use directive phrasing such as " "'we recommend' or 'should be replaced'. Keep paragraphs factual and observation-only. " ) word_targets = ( "executive_summary ~45–85 words, property_description ~55–110 words, " "condition_assessment ~110–220 words, defects_and_risks ~55–120 words, " "recommendations ~30–70 words. " ) identity_block_clean = (identity_block or "").strip() identity_clause = ( "PROPERTY IDENTITY (NON-NEGOTIABLE, must be reused verbatim — do NOT introduce any " "different address, postcode, property type, surveyor name, or company name):\n" f"{identity_block_clean}\n\n" if identity_block_clean else "" ) return ( "You are a Chartered Building Surveyor (MRICS) producing a RICS-style inspection report section. " "You decide which retrieval tools to call and how many times—there is no fixed script for RAG/KB. " "Ground factual statements in retrieve_survey_rag and/or retrieve_rics_kb results; " "use scan_note_duplicates when messy notes might duplicate other sections or prior library text. " "Never invent site-specific facts. If a fact is missing from bullets, seed evidence, retrieve_* " "results, and PROPERTY IDENTITY pins, omit that claim instead of using placeholder text. " "Never inline missing-data wording in the middle of a sentence. " + identity_clause + # ── Anti-summarisation + per-tier depth directive ───────────────── # The depth directive replaces the previous static word-target list, # which gave the same range for every level. With per-tier targets # gpt-4o-mini stops collapsing L3 sections into HomeBuyer-length # paragraphs. "PRODUCE COMPREHENSIVE PROPERTY-SPECIFIC ANALYSIS, NOT SUMMARIES. The five free-text fields " "in submit_inspection_section ARE the user-visible report. Each one must be substantive " "RICS-grade prose grounded in bullets and retrieved snippets — for THIS tier: " + word_targets + "Do NOT pad with generic boilerplate. Do NOT shorten to one or two sentences. " "When a topic is sparsely evidenced, retrieve more with retrieve_survey_rag or retrieve_rics_kb " "rather than truncating. Use the approved missing-info phrases ONLY for facts you genuinely " "cannot verify in the bullets/snippets — never as a substitute for analysis. Surface all " "evidenced defects from the bullets; do not omit any. " + depth_directive + f"Writing tone guidance: {tone}. " f"Server context (inject into tool calls automatically where applicable): tenant_id={tenant_id!r}, " f"primary_document_id={primary_document_id!r}. " f"RICS product tier: Level {pack.level} — {pack.product_label}. " f"Client survey_level field: {survey_level!r} (None is treated as Level 3 for template pack resolution). " f"Mandatory section order for this tier (codes): {order_preview}. " f"Current section: {section_code} — {section_title}. " "SERVER-ENFORCED WORKFLOW (tools): " "(1) retrieve evidence as needed → " "(2) submit_extraction_audit → " "(3) submit_section_plan → " f"(4){' submit_condition_rating_summary →' if section_code == 'C' and ratings_pack else ''} " "(5) submit_inspection_section (exactly once). " "Do not call submit_inspection_section in the same assistant message as any other tool. " "UK English." f"{c_extra}" ) async def run_inspector_tool_loop( *, db, tenant_id: str, primary_document_id: str, section_code: str, bullets: list[str], style_profile: WritingStyleProfile | None, ai_percent: int, retrieval_level: str, reference_document_ids: list[str] | None, peer_sections: dict[str, str] | None, survey_level: int | None = None, ) -> tuple[StructuredReport, list[dict[str, Any]]]: """Run the autonomous inspector loop; returns structured report + tool trace.""" trace: list[dict[str, Any]] = [] from app.agentic.speculative_executor import ( PatternRegistry, SpeculativeToolDispatcher, speculative_execution_active, ) _spec_registry = PatternRegistry() _spec_context: dict[str, Any] = { "section_code": section_code, "outline": "", } pack = get_survey_pack(survey_level) template = get_template(section_code, survey_level) section_title = template.title if template else section_code skeleton = (template.skeleton if template else "")[:2000] _spec_context["outline"] = skeleton # ── Pre-seeded grounding (key fix for "LLM invents because it had no evidence") # Without an upfront retrieval seed, the inspector loop relied entirely on # the LLM to discover useful queries via retrieve_survey_rag. When the # bullets are sparse (the user's exact complaint scenario — messy notes # uploaded for an L3 section), the LLM's first query is often vague # ("section E2 condition"), the result set is bland, the LLM gives up # retrieving and falls back to its training-data prior. That's where # "10 Kingsley Avenue", "robust steel frame", and made-up surveyor names # come from. # # We now do a deterministic up-front retrieval against the primary # uploaded survey (and any user-tagged reference docs) and inject the # top hits straight into the user message as SEED EVIDENCE, plus # extracted IDENTITY FACTS and STYLE REFERENCE PARAGRAPHS from the # tenant's profile. The LLM still has full freedom to call additional # retrieve_* tools, but starts with concrete grounded material — which # both reduces invention and gives the verifier real text to compare # the draft against. # # Local imports keep the cycle clean: the agent tools layer already # imports inspector_loop indirectly via agents.py. from app.agentic import tools as agent_tools_local from app.services.generation import _extract_property_identity, _identity_facts_block seed_hits: list[SearchResult] = [] seed_query = " ".join( [section_code, section_title] + [b for b in bullets if isinstance(b, str) and b.strip()][:12] ).strip() # Tier-aware seed retrieval window. Previously L1 and L2 received the # same window (k=16/rerank=6), which under-served L2 (HomeBuyer needs # enough evidence to give *proportionate* advice — i.e. not just # observation, but options-light) and over-served L1 (observation-only, # narrower scope). The new split: # L1 — k=12 / rerank=5 (observation-only, condensed RAG window) # L2 — k=18 / rerank=8 (HomeBuyer; more evidence than L1, less than L3) # L3 — k=24 / rerank=10 (Building Survey; full diagnostic context) effective_level = survey_level if survey_level is not None else pack.level try: _seed_lvl = int(effective_level) except Exception: # noqa: BLE001 _seed_lvl = pack.level if _seed_lvl >= 3: seed_k, seed_rerank = 24, 10 elif _seed_lvl == 2: seed_k, seed_rerank = 18, 8 else: seed_k, seed_rerank = 12, 5 if seed_query: try: seed_hits = await agent_tools_local.retrieve_tenant_evidence_async( query=seed_query, tenant_id=tenant_id, primary_document_id=primary_document_id, secondary_document_ids=list(reference_document_ids or []), k=seed_k, rerank_top_n=seed_rerank, ) except Exception as exc: # noqa: BLE001 logger.warning("Inspector seed retrieval failed (%s); proceeding without pre-seed", exc) seed_hits = [] trace.append({"event": "seed_retrieval", "hits": len(seed_hits), "query_preview": seed_query[:120]}) identity = _extract_property_identity(list(bullets)) if seed_hits: # Identity often appears in the primary doc but not in messy bullets # — e.g. the surveyor uploaded the property's own description PDF # but only typed a few free-form notes for THIS section. Letting the # extractor see the seed snippets too means the address it pins is # the address from the actual uploaded survey, not a parsed-out # bullet fragment. Trim each snippet to keep extraction fast. ext_lines = list(bullets) + [r.text[:1000] for r in seed_hits[:6] if getattr(r, "text", None)] identity = _extract_property_identity(ext_lines) identity_block = _identity_facts_block(identity) seed_evidence_block = "" if seed_hits: # Mirror the seed retrieval tier ladder: don't ship more snippets to # the LLM than were retrieved, and don't bloat the L1 prompt with # context the L1 mode won't use. if _seed_lvl >= 3: max_seed = 10 elif _seed_lvl == 2: max_seed = 8 else: max_seed = 5 seed_lines: list[str] = [] for i, r in enumerate(seed_hits[:max_seed], 1): t = (r.text or "").strip().replace("\n", " ") if len(t) > 600: t = t[:599] + "…" tag_kb = "(KB)" if bool(getattr(r, "kb", False)) else "(uploaded)" seed_lines.append(f"[#{i}] {tag_kb} {t}") seed_evidence_block = ( "\n\nSEED EVIDENCE FROM YOUR UPLOADED RAG CORPUS (primary grounding — facts in here " "are TRUSTED; quote or paraphrase rather than invent):\n" + "\n\n".join(seed_lines) ) style_examples_block = "" if style_profile is not None: examples = [ p for p in (getattr(style_profile, "example_paragraphs", []) or []) if isinstance(p, str) and p.strip() ][:3] if examples: style_examples_block = ( "\n\nSTYLE REFERENCE PARAGRAPHS (verbatim extracts from this surveyor's own completed " "RICS reports — mirror the tone, sentence rhythm, and phrasing patterns; do NOT copy " "facts from these unless the bullets/seed-evidence also support them):\n" + "\n\n".join(f'"{p.strip()[:600]}"' for p in examples) ) user_intro = ( f"RICS product pack: Level {pack.level} — {pack.product_label}\n" f"RICS section: {section_code} — {section_title}\n\n" f"Template skeleton (use as structural guide):\n{skeleton}\n\n" "Messy inspection notes (bullets — interpret, prioritise, and map to RICS discipline):\n" + "\n".join(f"- {b}" for b in bullets if str(b).strip()) + seed_evidence_block + style_examples_block ) if peer_sections: user_intro += "\n\nPeer section drafts (for duplicate scan if relevant):\n" + json.dumps( {str(k): str(v)[:800] for k, v in list(peer_sections.items())[:24]}, ensure_ascii=False, ) tools = _openai_inspector_tool_definitions() ai_params = _ai_level_to_params(3, ai_percent=ai_percent) temperature = float(ai_params["temperature"]) from app.llm.prompt_cache import ensure_cacheable_system_prefix, prompt_caching_active system_content = _inspector_system_prompt( tenant_id=tenant_id, primary_document_id=primary_document_id, survey_level=survey_level, pack=pack, section_code=section_code, section_title=section_title, style_profile=style_profile, identity_block=identity_block, ) if prompt_caching_active(): system_content = ensure_cacheable_system_prefix(system_content) messages: list[dict[str, Any]] = [ {"role": "system", "content": system_content}, {"role": "user", "content": user_intro}, ] # Seed hits also feed the eventual non-invention guard so it has concrete # comparison material. Without this, the verifier saw only the LLM-issued # retrieve_survey_rag output (which the LLM may have skipped). all_hits: list[SearchResult] = list(seed_hits) max_rounds = max(2, min(32, int(settings.inspector_max_tool_rounds))) submit_payload: dict[str, str] | None = None gate: dict[str, Any] = { "audit_payload": None, "plan_payload": None, "condition_payload": None, "submit_payload": None, } async def _dispatch_tool_wrapped(*, name: str, args: dict[str, Any]) -> tuple[str, list[SearchResult]]: return await _dispatch_tool( db=db, tenant_id=tenant_id, primary_document_id=primary_document_id, reference_document_ids=reference_document_ids, default_hierarchy_level=retrieval_level if retrieval_level in ("document", "section", "paragraph") else None, peer_sections_default=peer_sections, section_code=section_code, pack=pack, gate=gate, name=name, args=args, ) _spec_dispatcher: SpeculativeToolDispatcher | None = None if speculative_execution_active(): _spec_dispatcher = SpeculativeToolDispatcher( registry=_spec_registry, dispatch_fn=_dispatch_tool_wrapped, context=_spec_context, ) soft_deadline_round = max(2, int(max_rounds * 0.75)) for round_i in range(max_rounds): # Keep exploratory/tool-selection rounds on the cheaper chat_model. # Escalate to inspector_body_model only once the workflow is in the # drafting phase (audit + plan accepted) or when round budget is tight. drafting_phase = bool(gate.get("audit_payload")) and bool(gate.get("plan_payload")) current_model = ( settings.inspector_body_model if drafting_phase or round_i >= soft_deadline_round else settings.chat_model ) from app.llm.llm_throttle import make_cache_hit_slot, throttled_llm_call from app.llm.prompt_cache import ( log_openai_cache_usage, normalize_messages_for_caching, openai_extra_kwargs, prompt_caching_active, ) async_client = _inspector_async_client() max_tok = min(6144, max(2048, int(settings.max_output_tokens) * 6)) extra = ( openai_extra_kwargs( phase="inspector_tool_round", model=current_model, survey_level=survey_level, tenant_id=tenant_id, ) if prompt_caching_active() else {} ) cache_slot = make_cache_hit_slot() api_messages = ( normalize_messages_for_caching(messages) if prompt_caching_active() else messages ) async def _async_round() -> object: response = await async_client.chat.completions.create( model=current_model, messages=api_messages, tools=tools, tool_choice="auto", temperature=temperature, max_tokens=max_tok, **extra, ) if prompt_caching_active(): cache_slot[0] = log_openai_cache_usage( response, phase="inspector_tool_round", section_id=section_code, ) return response resp = await throttled_llm_call( phase="inspector_tool_round", section_id=section_code, cache_hit_out=cache_slot, call=_async_round, ) msg = resp.choices[0].message if not msg.tool_calls: messages.append({"role": "assistant", "content": msg.content or ""}) messages.append( { "role": "user", "content": ( "Continue with the SERVER-ENFORCED tool workflow: " "submit_extraction_audit → submit_section_plan → " f"{'submit_condition_rating_summary → ' if _section_c_requires_rating_summary(section_code=section_code, pack=pack) else ''}" "then submit_inspection_section (five fields). " "Use retrieve_survey_rag / retrieve_rics_kb / scan_note_duplicates as needed before the submits. " "If evidence is thin, do not invent facts; omit unsupported claims." ), } ) trace.append({"round": round_i, "note": "no_tool_calls_nudge", "model": current_model}) continue tcs = list(msg.tool_calls) names = [tc.function.name for tc in tcs] submit_count = sum(1 for n in names if n == "submit_inspection_section") if (submit_count and len(tcs) > 1) or submit_count > 1: err = ( "Invalid: submit_inspection_section must be the only tool call in the assistant message." if submit_count and len(tcs) > 1 else "Invalid: at most one submit_inspection_section tool call." ) trace.append( { "round": round_i, "error": "submit_mixed_with_other_tools" if submit_count and len(tcs) > 1 else "multiple_submits", } ) assistant_tool_calls = [ {"id": tc.id, "type": "function", "function": {"name": tc.function.name, "arguments": tc.function.arguments}} for tc in tcs ] messages.append({"role": "assistant", "content": msg.content or None, "tool_calls": assistant_tool_calls}) for tc in tcs: messages.append({"role": "tool", "tool_call_id": tc.id, "content": json.dumps({"ok": False, "error": err})}) continue assistant_tool_calls = [] for tc in tcs: assistant_tool_calls.append({"id": tc.id, "type": "function", "function": {"name": tc.function.name, "arguments": tc.function.arguments}}) messages.append( { "role": "assistant", "content": msg.content or None, "tool_calls": assistant_tool_calls, } ) for tc in tcs: name = tc.function.name raw = tc.function.arguments or "{}" try: args = json.loads(raw) except json.JSONDecodeError: args = {} trace.append({"round": round_i, "tool": name, "arguments_preview": raw[:400], "model": current_model}) if _spec_dispatcher is not None: body, hits = await _spec_dispatcher.dispatch(name=name, args=args) else: body, hits = await _dispatch_tool( db=db, tenant_id=tenant_id, primary_document_id=primary_document_id, reference_document_ids=reference_document_ids, default_hierarchy_level=retrieval_level if retrieval_level in ("document", "section", "paragraph") else None, peer_sections_default=peer_sections, section_code=section_code, pack=pack, gate=gate, name=name, args=args, ) all_hits.extend(hits) messages.append({"role": "tool", "tool_call_id": tc.id, "content": body[:24_000]}) if name == "submit_inspection_section": try: parsed = json.loads(body) except json.JSONDecodeError: parsed = {} if isinstance(parsed, dict) and parsed.get("ok") is True and gate.get("submit_payload"): submit_payload = gate["submit_payload"] break if submit_payload is not None: break if round_i >= soft_deadline_round and submit_payload is None: retrieve_calls = sum( 1 for t in trace if isinstance(t, dict) and str(t.get("tool", "")).startswith("retrieve_") ) if retrieve_calls >= 6: messages.append( { "role": "user", "content": ( "You have retrieved enough evidence. STOP retrieving. " "Call submit_extraction_audit, then submit_section_plan, then " "submit_inspection_section immediately using current evidence. " "Partial evidence is acceptable; no-submit fallback is worse." ), } ) trace.append({"round": round_i, "note": "force_submit_deadline_nudge"}) # Risk + compliance heuristics (same as scripted agent) for stable API shapes from app.agentic import agents as agents_mod risk_agent = agents_mod.RiskAssessmentAgent() comp_agent = agents_mod.StandardsComplianceAgent() # `RiskAssessmentAgent.assess` is async (it can call OpenAI for contextual risk # reasoning, falling back to keyword heuristics offline). Without the await, # `risks` would be a coroutine and the downstream `tuple(risks)` would raise # TypeError, crashing risk/compliance finalization for every section. risks = await risk_agent.assess(section_code=section_code, bullets=bullets) kb_glimpses = [r.text for r in all_hits if bool(getattr(r, "kb", False)) or r.tenant_id == settings.knowledge_base_tenant_id][:3] comp = comp_agent.check( section_code=section_code, bullets=bullets, kb_guidance=kb_glimpses or None, survey_level=survey_level, ) uniq_hits = agent_tools.dedupe_search_results(all_hits) evidence_items = await _evidence_items_async(db, tenant_id, uniq_hits) if submit_payload: # Non-invention enforcement on EVERY LLM-emitted artifact. # # The five body fields (the user-visible report) get the full 2-layer # guard (regex + contextual LLM grounding). The auxiliary fields the # tool gates produced (extraction_audit, section_plan, # condition_rating_summary) and LLM-drafted risks all go through the # cheap synchronous regex pass. Without this, GPT-4o-mini routinely # ships invented postcodes / surveyor names / firm references in the # metadata bundle the API surfaces (`/agentic/...` returns the gate # payloads in `inspector_meta`). verify_snippets = [r.text for r in uniq_hits if getattr(r, "text", None)] peer_values = [ v for v in (peer_sections or {}).values() if isinstance(v, str) and v.strip() ] full_bullets = list(bullets) + peer_values submit_payload = await _enforce_verify_submit_payload( submit_payload, bullets=bullets, snippets=verify_snippets, peer_sections=peer_sections, pinned_identity=identity, ) gate["audit_payload"] = _verify_audit_payload( gate.get("audit_payload"), bullets=full_bullets, snippets=verify_snippets, pinned_identity=identity, ) gate["plan_payload"] = _verify_plan_payload( gate.get("plan_payload"), bullets=full_bullets, snippets=verify_snippets, pinned_identity=identity, ) gate["condition_payload"] = _verify_condition_payload( gate.get("condition_payload"), bullets=full_bullets, snippets=verify_snippets, pinned_identity=identity, ) risks = _verify_risks( risks, bullets=full_bullets, snippets=verify_snippets, pinned_identity=identity, ) # ── Level-1 behavioural enforcement (parity with the legacy LCEL path) # The legacy `_run_generate` calls `_tier_validation_issues` and # retries once when an L1 draft contains advice phrasing. The # agentic inspector loop has no such retry — by the time we reach # this point the LLM has already submitted, and a second tool-loop # round-trip per section would double the cost. We instead apply a # deterministic regex sanitiser that strips advice sentences from # the five body fields and from L1 risk actions. This is the same # sanitiser the legacy path can fall back to when the retry still # returns advice, so behaviour matches across both pipelines. if _seed_lvl <= 1: submit_payload = strip_l1_advice_payload(submit_payload) risks = [ RiskItem( category=r.category, risk=r.risk, severity=r.severity, likelihood=r.likelihood, action=strip_l1_advice(r.action) if r.action else r.action, evidence=r.evidence, ) for r in risks ] rep = StructuredReport( executive_summary=_strip_missing_fact_phrase(submit_payload["executive_summary"] or ""), property_description=_strip_missing_fact_phrase(submit_payload["property_description"] or ""), condition_assessment=_strip_missing_fact_phrase(submit_payload["condition_assessment"] or ""), defects_and_risks=_strip_missing_fact_phrase(submit_payload["defects_and_risks"] or ""), recommendations=_strip_missing_fact_phrase(submit_payload["recommendations"] or ""), risks=tuple(risks), compliance=tuple(comp), evidence_items=tuple(evidence_items), tool_trace=tuple(trace), extraction_audit=gate.get("audit_payload") if isinstance(gate.get("audit_payload"), dict) else None, section_plan=gate.get("plan_payload") if isinstance(gate.get("plan_payload"), dict) else None, condition_rating_summary=gate.get("condition_payload") if isinstance(gate.get("condition_payload"), dict) else None, ) return rep, trace # Fallback: single LCEL draft from accumulated evidence (no submit) from app.llm import generation_facade as gen_llm snippets = [r.text for r in uniq_hits[:12]] base = await gen_llm.generate_section( skeleton=template.skeleton if template else f"[{section_code}]: [content].", bullets=bullets, snippets=[], style_profile=style_profile if ai_percent > 5 else None, temperature=temperature, creativity_hint=str(ai_params["creativity_hint"]) + "\n\nAGENTIC FALLBACK: Tool loop did not call submit_inspection_section in time; draft from evidence.", document_context=snippets[:4], hierarchy_section_snippets=None, paragraph_snippets=snippets[4:], tenant_id=tenant_id, style_anchor=None, survey_level=survey_level, ) # Same non-invention enforcement as the submit path. Fallback drafts skip # the structured five-field tool, so they are arguably MORE prone to # hallucination — guarding here is critical, not optional. try: base = await async_enforce_verify( text=base, bullets=bullets, snippets=snippets, pinned_identity=identity, openai_api_key=settings.openai_api_key or "", model=settings.chat_model, ) except Exception as exc: # noqa: BLE001 logger.warning("Non-invention guard failed on fallback draft (%s); returning raw text", exc) trace.append({"error": "submit_inspection_section not called", "fallback": "generate_section"}) # Apply the same regex guard to fallback risks for parity with the success # path. (`base` itself was already passed through async_enforce_verify # above; this catches inventions inside risk[*].risk / risk[*].action.) risks = _verify_risks( risks, bullets=bullets, snippets=snippets, pinned_identity=identity, ) # Level-1 behavioural enforcement on the fallback draft and on the # risk[*].action strings — mirrors the success path's L1 sweep above. if _seed_lvl <= 1: base = strip_l1_advice(base) risks = [ RiskItem( category=r.category, risk=r.risk, severity=r.severity, likelihood=r.likelihood, action=strip_l1_advice(r.action) if r.action else r.action, evidence=r.evidence, ) for r in risks ] # Tier-aware defects/recommendations construction. Building both fields # by concatenating risk[*].risk + risk[*].action is the right shape for # L2/L3 (HomeBuyer / Building Survey both produce explicit # recommendations), but writes "Investigate source" / "Seek structural # engineer review" into an L1 product — exactly the leakage the L1 # sanitiser above is meant to prevent. For L1 we surface defects as # observation only and emit the L1 placeholder for recommendations. if _seed_lvl <= 1: defects_text = " ".join( f"{r.category}: {r.risk} ({r.severity}/{r.likelihood})." for r in risks ).strip()[:2400] recommendations_text = _L1_PLACEHOLDER else: defects_text = " ".join( f"{r.category}: {r.risk} ({r.severity}/{r.likelihood}). {r.action}" for r in risks ).strip()[:2400] recommendations_text = " ".join(r.action for r in risks).strip()[:1600] # Previous behaviour duplicated `base[:900]` across both # `property_description` and `condition_assessment`, and used `base[:400]` # as a third copy in `executive_summary` — three views of the same # truncated text, which is exactly why fallback output read like a # one-paragraph summary repeated under three headings. We instead place # the full LLM draft once, in `condition_assessment` (the spine of the # L3 element renderer and the most substantive subhead for non-L3 codes # in `render_report_text`), and leave `property_description` blank so # the renderer skips it rather than printing a duplicate. # Fallback path: the autonomous tool loop did not converge on a clean # structured submit within max_rounds. We still have a usable LLM # narrative draft (`base`) collected from retrieve_survey_rag snippets, # so we hand THAT to the user under condition_assessment. The previous # version of this branch leaked dev-quality language ("the autonomous # tool loop did not complete a structured submit before the round # limit") into executive_summary — surfacing internal pipeline state in # the user-facing report. RAGAS evaluation flagged this as the largest # qualitative defect on Section D in the 37 Elms Crescent run. The # replacement copy below is short, professional, and indistinguishable # from a clean submit at the user level; the underlying telemetry is # still visible to operators via tool_trace and the pipeline metadata. rep = StructuredReport( executive_summary=( f"{section_title}: this section summarises the surveyor's notes and the " f"evidence drawn from the supporting documents. Key observations and any " f"matters identified during the inspection are set out below." ).strip(), property_description="", condition_assessment=_strip_missing_fact_phrase(base.strip()), defects_and_risks=_strip_missing_fact_phrase(defects_text), recommendations=_strip_missing_fact_phrase(recommendations_text), risks=tuple(risks), compliance=tuple(comp), evidence_items=tuple(evidence_items), tool_trace=tuple(trace), extraction_audit=None, section_plan=None, condition_rating_summary=None, ) return rep, trace