RandomZ / app /agentic /inspector_loop.py
StormShadow308's picture
Ship personalised RAG, 10m SLA, parallel generate, and AI phase docs for HF pilot.
be9fd4a
Raw
History Blame Contribute Delete
78.8 kB
"""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