letxinet / backend /synthesis.py
C2MV's picture
Initial upload for Build Small Hackathon
68fb5e2 verified
Raw
History Blame Contribute Delete
40.1 kB
"""
Synthesis Engine - AI-powered research synthesis pipeline
Enhanced with hierarchical synthesis, GRADE classification, gap detection, and rescue search.
Faithful to the original Next.js research-agent prompts.
"""
import json
import httpx
import re
from typing import Dict, Any, List, Optional
from .prompts.profiles import AGENT_PROFILES
from .prompts.synthesis import (
MASTER_SYNTHESIS_PROMPT,
WRITING_PROMPT,
VALIDATION_PROMPT,
AUDIT_PROMPT,
ARA_PROMPT,
)
from .prompts.planning import SEARCH_PLANNING_PROMPT, GAP_DETECTION_PROMPT
from .utils import robust_json_parse, extract_research_plan
# Provider configs
PROVIDERS = {
"groq": {
"base_url": "https://api.groq.com/openai/v1",
"env_key": "GROQ_API_KEY",
"models": [
"llama-3.3-70b-versatile",
"llama-3.1-8b-instant",
"deepseek-r1-distill-llama-70b",
"mixtral-8x7b-32768",
"gemma2-9b-it",
"llama3-70b-8192",
"llama3-8b-8192",
"llama-guard-3-8b",
],
},
"openrouter": {
"base_url": "https://openrouter.ai/api/v1",
"env_key": "OPENROUTER_API_KEY",
"models": [
"meta-llama/llama-3.3-70b-instruct:free",
"google/gemma-4-26b-a4b-it:free",
"google/gemma-4-31b-it:free",
"nvidia/nemotron-3-super-120b-a12b:free",
"deepseek/deepseek-v4-flash:free",
"deepseek/deepseek-r1-0528:free",
"qwen/qwen3-next-80b-a3b-instruct:free",
"minimax/minimax-m2.5:free",
"openai/gpt-oss-120b:free",
"openai/gpt-oss-20b:free",
"arcee-ai/trinity-large-thinking:free",
"nousresearch/hermes-3-llama-3.1-405b:free",
"google/gemma-3-27b-it:free",
"google/gemma-3-12b-it:free",
"qwen/qwen3-coder:free",
"stepfun/step-3.5-flash:free",
"z-ai/glm-4.5-air:free",
"anthropic/claude-sonnet-4.5",
"anthropic/claude-haiku-4.5",
"openai/gpt-5.4",
"openai/gpt-5.4-mini",
"openai/gpt-5",
"deepseek/deepseek-v4-pro",
"deepseek/deepseek-v3.2",
"qwen/qwen3.6-flash",
"qwen/qwen3.5-plus-20260420",
"mistralai/mistral-small-2603",
"mistralai/mistral-medium-3-5",
],
},
"mistral": {
"base_url": "https://api.mistral.ai/v1",
"env_key": "MISTRAL_API_KEY",
"models": [
"mistral-small-2506",
"mistral-small-2603",
"mistral-medium-2508",
"mistral-medium-3-5",
"mistral-large-2512",
"magistral-medium-2509",
"magistral-small-2509",
"ministral-3b-2512",
"ministral-8b-2512",
"ministral-14b-2512",
"codestral-2508",
"devstral-2512",
"open-mistral-nemo",
],
},
"gemini": {
"base_url": "https://generativelanguage.googleapis.com/v1beta/openai",
"env_key": "GEMINI_API_KEY",
"models": [
"gemini-2.5-flash",
"gemini-2.5-pro",
"gemini-2.0-flash",
"gemini-2.0-flash-lite",
"gemini-3-flash-preview",
"gemini-3-pro-preview",
"gemini-3.1-flash-lite",
"gemma-4-26b-a4b-it",
"gemma-4-31b-it",
],
},
"deepseek": {
"base_url": "https://api.deepseek.com/v1",
"env_key": "DEEPSEEK_API_KEY",
"models": [
"deepseek-chat",
"deepseek-reasoner",
"deepseek-v4-flash",
"deepseek-v4-pro",
],
},
"nebius": {
"base_url": "https://api.tokenfactory.nebius.com/v1",
"env_key": "NEBIUS_API_KEY",
"models": [
"deepseek-ai/DeepSeek-V3.2",
"deepseek-ai/DeepSeek-V4-Pro",
"meta-llama/Llama-3.3-70B-Instruct",
"Qwen/Qwen3-235B-A22B-Instruct-2507",
"Qwen/Qwen3-32B",
"Qwen/Qwen3.5-397B-A17B",
"nvidia/Llama-3_1-Nemotron-Ultra-253B-v1",
"google/gemma-3-27b-it",
"NousResearch/Hermes-4-405B",
"moonshotai/Kimi-K2.5",
"MiniMaxAI/MiniMax-M2.5",
],
},
"azure": {
"base_url": "https://letxinet.openai.azure.com/openai/deployments",
"env_key": "AZURE_API_KEY",
"models": [
"gpt-4o-mini",
"gpt-4o",
"o3-mini",
"o4-mini",
"gpt-4.1-mini",
],
},
"huggingface": {
"base_url": "https://api-inference.huggingface.co/v1",
"env_key": "HF_TOKEN",
"models": [
"deepseek-ai/DeepSeek-V3.2",
"deepseek-ai/DeepSeek-R1",
"meta-llama/Llama-3.3-70B-Instruct",
"meta-llama/Llama-4-Scout-17B-16E-Instruct",
"Qwen/Qwen3-235B-A22B-Instruct-2507",
"Qwen/Qwen3-Next-80B-A3B-Instruct",
"google/gemma-3-27b-it",
"MiniMaxAI/MiniMax-M2.1",
"moonshotai/Kimi-K2.5",
],
},
}
GRADE_LEVELS = {
"1a": {"label": "Meta-análisis", "weight": 10, "desc": "Revisión sistemática cuantitativa con pooling estadístico"},
"1b": {"label": "Revisión sistemática", "weight": 9, "desc": "Búsqueda exhaustiva y replicable con criterios de inclusión/exclusión"},
"2a": {"label": "Ensayo controlado aleatorizado", "weight": 8, "desc": "Experimento con aleatorización y grupo control"},
"2b": {"label": "Ensayo cuasi-experimental", "weight": 7, "desc": "Experimento sin aleatorización completa"},
"3a": {"label": "Estudio de cohorte", "weight": 6, "desc": "Seguimiento longitudinal de grupos expuestos/no expuestos"},
"3b": {"label": "Estudio caso-control", "weight": 5, "desc": "Comparación retrospectiva de casos y controles"},
"4": {"label": "Corte transversal", "weight": 4, "desc": "Medición en un punto único del tiempo"},
"5": {"label": "Serie de casos", "weight": 3, "desc": "Descripción de grupos sin grupo control"},
"6": {"label": "Opinión de expertos", "weight": 2, "desc": "Juicio clínico o consenso de especialistas"},
}
OXFORD_LEVELS = {
"1a": {"label": "RS de ensayos aleatorizados", "weight": 10, "desc": "Revisión Sistemática de RCTs"},
"1b": {"label": "Ensayo controlado aleatorizado", "weight": 9, "desc": "RCT individual con intervalo de confianza estrecho"},
"1c": {"label": "Todo o nada", "weight": 8, "desc": "Todos los pacientes murieron antes que estuviera disponible el tratamiento, y ahora algunos sobreviven; o cuando algunos pacientes morían antes de que estuviera disponible el tratamiento, y ahora ninguno muere"},
"2a": {"label": "RS de estudios de cohorte", "weight": 7, "desc": "Revisión Sistemática de estudios de cohorte"},
"2b": {"label": "Estudio de cohorte", "weight": 6, "desc": "Estudio de cohorte individual o RCT de baja calidad"},
"2c": {"label": "Investigación de resultados", "weight": 5, "desc": "Investigación de resultados, estudios ecológicos"},
"3a": {"label": "RS de estudios caso-control", "weight": 4, "desc": "Revisión Sistemática de estudios caso-control"},
"3b": {"label": "Estudio caso-control", "weight": 3, "desc": "Estudio caso-control individual"},
"4": {"label": "Serie de casos", "weight": 2, "desc": "Serie de casos, o estudios de cohorte o de caso-control de baja calidad"},
"5": {"label": "Opinión de expertos", "weight": 1, "desc": "Opinión de expertos sin evaluación crítica explícita"},
}
ORIGINAL_GRADE_LEVELS = {
"ALTA": {
"label": "ALTA",
"weight": 4,
"desc": "Meta-analisis, revisiones sistematicas o ensayos controlados aleatorizados.",
},
"MODERADA": {
"label": "MODERADA",
"weight": 3,
"desc": "Ensayos clinicos, estudios experimentales controlados, cohortes o casos y controles bien disenados.",
},
"BAJA": {
"label": "BAJA",
"weight": 2,
"desc": "Estudios observacionales, descriptivos o transversales.",
},
"MUY BAJA": {
"label": "MUY BAJA",
"weight": 1,
"desc": "Reportes de caso, opiniones, editoriales o evidencia no revisada.",
},
}
ORIGINAL_GRADE_ALIASES = {
"ALTO": "ALTA",
"HIGH": "ALTA",
"ALTA": "ALTA",
"MODERADO": "MODERADA",
"MODERADA": "MODERADA",
"MODERATE": "MODERADA",
"MEDIUM": "MODERADA",
"BAJO": "BAJA",
"BAJA": "BAJA",
"LOW": "BAJA",
"MUY BAJO": "MUY BAJA",
"MUY BAJA": "MUY BAJA",
"VERY LOW": "MUY BAJA",
"VERY_LOW": "MUY BAJA",
}
def normalize_original_grade_level(level: Any) -> str:
"""Normalize original beta GRADE labels to ALTA/MODERADA/BAJA/MUY BAJA."""
raw = str(level or "").strip().upper().replace("_", " ")
raw = re.sub(r"\s+", " ", raw)
return ORIGINAL_GRADE_ALIASES.get(raw, "BAJA")
def classify_grade_original(study_type: str) -> str:
"""Fast fallback that maps study design keywords to the original beta GRADE labels."""
numeric = classify_grade(study_type)
if numeric in {"1a", "1b", "2a"}:
return "ALTA"
if numeric in {"2b", "3a", "3b"}:
return "MODERADA"
if numeric in {"4", "5"}:
return "BAJA"
return "MUY BAJA"
def def_document_has_grade(doc: Dict[str, Any]) -> bool:
return bool(doc.get("grade_level") or doc.get("evidenceLevel"))
def classify_grade_oxford(study_type: str) -> str:
"""Classify a study type string into Oxford CEBM evidence level."""
t = study_type.lower()
if "revisión sistemática" in t and ("aleatorizado" in t or "rct" in t):
return "1a"
if "meta-análisis" in t or "meta-analisis" in t or "meta analysis" in t:
return "1a"
if "ensayo" in t and ("aleatorizado" in t or "randomized" in t or "rct" in t):
return "1b"
if "revisión sistemática" in t and ("cohorte" in t or "cohort" in t):
return "2a"
if "cohorte" in t or "cohort" in t or "longitudinal" in t:
return "2b"
if "ecológico" in t or "ecological" in t:
return "2c"
if "revisión sistemática" in t and ("caso-control" in t or "case-control" in t):
return "3a"
if "caso-control" in t or "case-control" in t:
return "3b"
if "serie de casos" in t or "case series" in t or "transversal" in t or "cross-sectional" in t or "encuesta" in t:
return "4"
if "experto" in t or "opinión" in t or "expert" in t:
return "5"
return "4" # Default
def classify_grade(study_type: str) -> str:
"""Classify a study type string into GRADE evidence level."""
t = study_type.lower()
if "meta-análisis" in t or "meta-analisis" in t or "meta analysis" in t:
return "1a"
if "revisión sistemática" in t or "revision sistematica" in t or "systematic review" in t:
return "1b"
if "ensayo" in t and ("aleatorizado" in t or "randomized" in t or "rct" in t):
return "2a"
if "ensayo" in t or "quasi" in t or "quasi-experimental" in t:
return "2b"
if "cohorte" in t or "cohort" in t or "longitudinal" in t:
return "3a"
if "caso-control" in t or "case-control" in t:
return "3b"
if "transversal" in t or "cross-sectional" in t or "encuesta" in t:
return "4"
if "serie de casos" in t or "case series" in t:
return "5"
if "experto" in t or "opinión" in t or "expert" in t:
return "6"
return "4"
def grade_label(level: str) -> str:
entry = GRADE_LEVELS.get(level, GRADE_LEVELS["4"])
return f"[{level.upper()}] {entry['label']}"
def grade_weight(level: str) -> int:
return GRADE_LEVELS.get(level, GRADE_LEVELS["4"])["weight"]
class SynthesisEngine:
def __init__(self, provider: str = "mistral", model: str = None, api_key: str = None,
search_model: str = None, translation_model: str = None):
config = PROVIDERS.get(provider, PROVIDERS["mistral"])
self.base_url = config["base_url"]
self.model = model or "mistral-small-2506"
self.search_model = search_model or self.model
self.translation_model = translation_model or self.model
self.api_key = api_key or ""
self.client = httpx.AsyncClient(timeout=180.0)
async def _call_llm(self, system_prompt: str, user_prompt: str, temperature: float = 0.0, role: str = "synthesis") -> str:
# Select model based on role
model_map = {
"search": self.search_model,
"synthesis": self.model,
"translation": self.translation_model,
}
active_model = model_map.get(role, self.model)
headers = {
"Authorization": f"Bearer {self.api_key}",
"Content-Type": "application/json",
}
payload = {
"model": active_model,
"messages": [
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_prompt},
],
"temperature": temperature,
"max_tokens": 8192,
}
# Clamp max_tokens based on model capabilities
MODEL_MAX_TOKENS = {
"mistral": 8192, "groq": 32768, "openrouter": 8192,
"gemini": 65536, "deepseek": 8192, "nebius": 8192,
"azure": 16384, "huggingface": 8192,
}
provider_key = self.base_url.split("//")[-1].split(".")[0] if "//" in self.base_url else "default"
max_allowed = MODEL_MAX_TOKENS.get(provider_key, 8192)
# Special cap for DeepSeek models
if "deepseek" in active_model.lower():
max_allowed = min(max_allowed, 8192)
requested_tokens = min(payload.get("max_tokens", 8192), max_allowed)
payload["max_tokens"] = requested_tokens
try:
r = await self.client.post(
f"{self.base_url}/chat/completions", json=payload, headers=headers
)
r.raise_for_status()
return r.json()["choices"][0]["message"]["content"]
except Exception as e:
raise RuntimeError(f"Error calling LLM: {str(e)}") from e
def _parse_json(self, text: str) -> dict:
result = robust_json_parse(text)
if result is not None:
return result
return {"error": "Could not parse JSON", "raw": text[:500]}
# ── Core pipeline phases (faithful to original prompts) ──────────────
async def orchestrate(self, query: str) -> dict:
"""Phase 1: Analyze query and extract variables."""
system = "Eres un orquestador de investigación académica. Analiza la consulta y extrae variables."
user = f"""Analiza esta consulta de investigación y extrae:
1. Sujeto de estudio
2. Variable Independiente (V.I.) con dimensiones
3. Variable Dependiente (V.D.) con dimensiones
4. Tipo de estudio sugerido
5. País/Contexto geográfico
6. Keywords en español e inglés
CONSULTA: "{query}"
RESPONDE EN JSON:
{{
"subject": "...",
"variable_independiente": {{"nombre": "...", "dimensiones": [...], "indicadores": [...]}},
"variable_dependiente": {{"nombre": "...", "dimensiones": [...], "indicadores": [...]}},
"tipo_estudio": "...",
"country": "...",
"keywords_es": [...],
"keywords_en": [...]
}}"""
response = await self._call_llm(system, user, role="search")
return self._parse_json(response)
async def plan_search(self, query: str, profile: str = "general", orchestrator_ctx: dict = None) -> dict:
"""Phase 2: Plan search queries."""
profile_data = AGENT_PROFILES.get(profile, AGENT_PROFILES["general"])
system = f"Eres un estratega de búsqueda académica. {profile_data['title']}."
user = SEARCH_PLANNING_PROMPT.format(query=query, agent_role=profile)
response = await self._call_llm(system, user, role="search")
return self._parse_json(response)
async def generate_master_plan(
self,
query: str,
docs_context: str,
profile: str = "general",
template_structure: str = None,
geo_context: str = "Automático",
) -> dict:
"""Phase 3: Generate master synthesis plan (linear path)."""
profile_data = AGENT_PROFILES.get(profile, AGENT_PROFILES["general"])
system = f"Eres un {profile_data['title']}. Genera un plan maestro de investigación. Contexto Geográfico: {geo_context}"
user = MASTER_SYNTHESIS_PROMPT.format(
query=query,
agent_title=profile_data["title"],
agent_title_upper=profile_data["title"].upper(),
profile_instruction=profile_data["instruction"],
template_structure=template_structure or "Genera la estructura que consideres adecuada.",
)
user += f"\n\nCONTEXTO GEOGRÁFICO ASIGNADO: {geo_context}"
user += f"\n\nDOCUMENTOS ENCONTRADOS:\n{docs_context}"
response = await self._call_llm(system, user, temperature=0.0, role="synthesis")
return extract_research_plan(response)
async def write_section(self, section_name: str, section_prompt: str, context_text: str, geo_context: str = "Automático") -> str:
"""Phase 4: Write individual section content."""
system = f"Eres un Redactor Científico Experto. Contexto Geográfico a priorizar: {geo_context}"
user = WRITING_PROMPT.replace("{section}", section_name).replace("{section_prompt}", section_prompt).replace("{context_text}", context_text)
user += f"\n\nCONTEXTO GEOGRÁFICO A PRIORIZAR: {geo_context}"
return await self._call_llm(system, user, temperature=0.0, role="synthesis")
async def validate_citations(self, docs_context: str, content: str) -> dict:
"""Phase 5a: Validate citations."""
system = "Eres un Agente de Validación Bibliográfica ESTRICTO."
user = VALIDATION_PROMPT.replace("{docs_context}", docs_context).replace("{content_to_validate}", content)
response = await self._call_llm(system, user, temperature=0.0, role="synthesis")
return self._parse_json(response)
async def audit_content(self, docs_context: str, content: str) -> dict:
"""Phase 5b: Audit content quality."""
system = "Eres un Auditor Técnico de Calidad Académica."
user = AUDIT_PROMPT.replace("{docs_context}", docs_context).replace("{content_to_audit}", content)
response = await self._call_llm(system, user, temperature=0.0, role="synthesis")
return self._parse_json(response)
async def refine_section(self, section_content: str, findings: str) -> str:
"""Phase 5c: Refine section with ARA+."""
system = "Eres el Agente de Refinamiento Académico Avanzado (ARA+)."
user = ARA_PROMPT.replace("{section_content}", section_content).replace("{section_findings}", findings)
return await self._call_llm(system, user, temperature=0.0, role="synthesis")
async def detect_gaps(self, query: str, plan_sections: list) -> dict:
"""Detect gaps in the research plan."""
system = "Eres un Auditor de Cobertura Científica."
sections_json = json.dumps(plan_sections)
user = GAP_DETECTION_PROMPT.replace("{query}", query).replace("{plan_sections}", sections_json)
response = await self._call_llm(system, user, temperature=0.0, role="search")
return self._parse_json(response)
# ── GRADE evidence classification ───────────────────────────────────
def _enrich_with_grade(
self,
doc: Dict[str, Any],
level: str,
system: str = "grade",
evidence_type: str = "",
justification: str = "",
) -> Dict[str, Any]:
"""Helper to attach grade metadata to a document based on system ('grade' or 'oxford')."""
if system == "original":
normalized = normalize_original_grade_level(level)
entry = ORIGINAL_GRADE_LEVELS[normalized]
return {
**doc,
"grade_level": normalized,
"grade_label": entry["label"],
"grade_weight": entry["weight"],
"grade_desc": entry["desc"],
"grade_system": "original",
"evidenceLevel": entry["label"],
"type": evidence_type or doc.get("type") or doc.get("study_type") or "",
"grade_justification": justification or doc.get("grade_justification", ""),
}
if system == "oxford":
entry = OXFORD_LEVELS.get(level, OXFORD_LEVELS["4"])
label = f"[{level.upper()}] {entry['label']}"
weight = entry["weight"]
desc = entry["desc"]
else:
entry = GRADE_LEVELS.get(level, GRADE_LEVELS["4"])
label = f"[{level.upper()}] {entry['label']}"
weight = entry["weight"]
desc = entry["desc"]
return {
**doc,
"grade_level": level,
"grade_label": label,
"grade_weight": weight,
"grade_desc": desc,
"grade_system": system,
"evidenceLevel": doc.get("evidenceLevel") or label,
}
def _extract_grade_classifications(self, parsed: Any) -> List[Dict[str, Any]]:
"""Recover classifications from the original beta response shape and common variants."""
if isinstance(parsed, list):
return [x for x in parsed if isinstance(x, dict)]
if not isinstance(parsed, dict):
return []
for key in (
"classifications",
"grades",
"results",
"documents",
"analysis",
"items",
"data",
"evaluations",
"plan",
):
value = parsed.get(key)
if isinstance(value, list):
return [x for x in value if isinstance(x, dict)]
if isinstance(value, dict):
return [value]
if any(k in parsed for k in ("index", "level", "type")):
return [parsed]
return []
def _grade_docs_context(self, documents: List[Dict[str, Any]], limit: int) -> str:
lines = []
for i, doc in enumerate(documents[:limit], 1):
authors = doc.get("authors", [])
if isinstance(authors, list):
authors = ", ".join(str(a) for a in authors if a)
snippet = doc.get("abstract") or doc.get("snippet") or doc.get("summary") or ""
lines.append(
f"[{i}] Titulo: {doc.get('title', 'Sin titulo')} | "
f"Autores: {authors or 'No especificados'} | "
f"Resumen: {str(snippet)[:500]}"
)
return "\n\n".join(lines)
async def classify_documents(self, documents: List[Dict[str, Any]], mode: str = "keywords") -> List[Dict[str, Any]]:
"""
Classify documents using the specified strategy.
Modes: 'keywords' (default, fast), 'llm' (accurate but slow), 'oxford' (CEBM fast), 'hybrid' (keywords + llm for unknown).
"""
from backend.prompts.synthesis import GRADE_PROMPT, GRADE_ORIGINAL_PROMPT
import json
enriched = []
mode = mode.lower()
if mode == "original":
limit = min(50, len(documents))
parsed: Any = {}
try:
user = GRADE_ORIGINAL_PROMPT.format(
documents_text=self._grade_docs_context(documents, limit)
)
response = await self._call_llm(
"Eres un Agente de Evaluacion Metodologica. Tu salida debe ser exclusivamente JSON valido.",
user,
temperature=0.0,
role="synthesis",
)
parsed = self._parse_json(response)
except Exception as e:
print(f"[GRADE ORIGINAL] Error classifying docs: {e}. Falling back to keywords.")
classifications = self._extract_grade_classifications(parsed)
by_index = {}
for i, item in enumerate(classifications, 1):
try:
idx = int(item.get("index", i)) - 1
except (TypeError, ValueError):
idx = i - 1
by_index[idx] = item
for i, doc in enumerate(documents):
item = by_index.get(i)
if item and i < limit:
enriched.append(
self._enrich_with_grade(
doc,
item.get("level", "BAJA"),
"original",
evidence_type=item.get("type", ""),
justification=item.get("justification", item.get("reason", "")),
)
)
else:
study_type = doc.get("study_type", doc.get("type", "transversal"))
enriched.append(
self._enrich_with_grade(
doc,
classify_grade_original(study_type),
"original",
evidence_type=study_type,
)
)
return enriched
if mode == "keywords":
for doc in documents:
study_type = doc.get("study_type", doc.get("type", "transversal"))
level = classify_grade(study_type)
enriched.append(self._enrich_with_grade(doc, level, "grade"))
elif mode == "oxford":
for doc in documents:
study_type = doc.get("study_type", doc.get("type", "transversal"))
level = classify_grade_oxford(study_type)
enriched.append(self._enrich_with_grade(doc, level, "oxford"))
elif mode in ["llm", "hybrid"]:
# For hybrid, pre-filter with keywords to save tokens
docs_to_llm = []
if mode == "hybrid":
for doc in documents:
study_type = doc.get("study_type", doc.get("type", "transversal"))
level = classify_grade(study_type)
if level != "4": # Confident classification
enriched.append(self._enrich_with_grade(doc, level, "grade"))
else:
docs_to_llm.append(doc)
else:
docs_to_llm = documents
if docs_to_llm:
# Prepare content for LLM (up to 30 docs to avoid context window issues)
limit = 30
content_to_grade = ""
for i, doc in enumerate(docs_to_llm[:limit]):
authors_str = ", ".join(doc.get("authors", []))
snippet = doc.get("snippet", doc.get("abstract", ""))
content_to_grade += f"[{i+1}] ID: {doc.get('id', i)} | Autores: {authors_str} | Resumen: {snippet}\n\n"
system = "Eres un experto en clasificación de evidencia científica y medicina basada en evidencia."
user = GRADE_PROMPT.format(documents_text=content_to_grade)
try:
response = await self._call_llm(system, user, temperature=0.1, role="synthesis")
results = self._parse_json(response)
if not isinstance(results, list):
if isinstance(results, dict) and "classifications" in results:
results = results["classifications"]
else:
results = [results]
# Map results back to documents
for i, doc in enumerate(docs_to_llm):
if i < limit and i < len(results):
res = results[i]
# Handle different response structures
level = res.get("level", "4")
if not level: level = "4"
enriched.append(self._enrich_with_grade(doc, level, "grade"))
else:
# Fallback for remaining docs
enriched.append(self._enrich_with_grade(doc, "4", "grade"))
except Exception as e:
print(f"[GRADE LLM] Error classifying docs: {e}. Falling back to keywords.")
for doc in docs_to_llm:
study_type = doc.get("study_type", doc.get("type", "transversal"))
level = classify_grade(study_type)
enriched.append(self._enrich_with_grade(doc, level, "grade"))
else:
# Fallback
for doc in documents:
study_type = doc.get("study_type", doc.get("type", "transversal"))
level = classify_grade(study_type)
enriched.append(self._enrich_with_grade(doc, level, "grade"))
return enriched
def classify_evidence(self, documents: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
"""Legacy synchronous wrapper. Use await classify_documents() instead."""
if any(def_document_has_grade(doc) for doc in documents):
return documents
import asyncio
try:
loop = asyncio.get_event_loop()
return loop.run_until_complete(self.classify_documents(documents, "keywords"))
except RuntimeError:
# If event loop is already running, we have to fall back to simple keyword matching
enriched = []
for doc in documents:
study_type = doc.get("study_type", doc.get("type", "transversal"))
level = classify_grade(study_type)
enriched.append(self._enrich_with_grade(doc, level, "grade"))
return enriched
def sort_by_evidence(self, documents: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
"""Sort documents by GRADE weight descending (strongest evidence first)."""
return sorted(documents, key=lambda d: d.get("grade_weight", 0), reverse=True)
def evidence_summary(self, documents: List[Dict[str, Any]]) -> Dict[str, Any]:
"""Produce a GRADE distribution summary."""
counts: Dict[str, int] = {}
for doc in documents:
lvl = doc.get("grade_level", "4")
counts[lvl] = counts.get(lvl, 0) + 1
if any(
doc.get("grade_system") == "original" or doc.get("grade_level") in ORIGINAL_GRADE_LEVELS
for doc in documents
):
return {
"distribution": [
{"level": l, "label": ORIGINAL_GRADE_LEVELS[l]["label"], "count": counts.get(l, 0)}
for l in ["ALTA", "MODERADA", "BAJA", "MUY BAJA"] if counts.get(l, 0) > 0
],
"total": len(documents),
}
levels_desc = ["1a", "1b", "2a", "2b", "3a", "3b", "4", "5", "6"]
return {
"distribution": [
{"level": l, "label": GRADE_LEVELS[l]["label"], "count": counts.get(l, 0)}
for l in levels_desc if counts.get(l, 0) > 0
],
"total": len(documents),
}
# ── Full-text retrieval helpers ─────────────────────────────────────
def extract_full_text(self, doc: Dict[str, Any]) -> str:
"""Extract the best available full text from a document entry."""
for key in ("full_text", "text", "content", "body", "extracted_text"):
val = doc.get(key)
if val and isinstance(val, str) and len(val.strip()) > 50:
return val.strip()
abstract = doc.get("abstract", doc.get("summary", ""))
if abstract:
return f"[Solo disponible resumen/abstract]\n{abstract.strip()}"
return "[No se encontró texto completo ni abstract para este documento]"
def build_full_text_context(self, documents: List[Dict[str, Any]], max_chars: int = 120000) -> str:
"""Build a concatenated full-text context from documents respecting char limit."""
sorted_docs = self.sort_by_evidence(documents)
parts: List[str] = []
total = 0
for i, doc in enumerate(sorted_docs, 1):
ref_id = doc.get("id", i)
title = doc.get("title", "Sin título")
authors = doc.get("authors", "Autor desconocido")
year = doc.get("year", "?")
grade = doc.get("evidenceLevel") or doc.get("grade_label", "")
text = self.extract_full_text(doc)
header = f"[{i}] (BIB:{ref_id}) {title} - {authors} ({year}) [{grade}]"
chunk = f"{header}\n{text}\n"
if total + len(chunk) > max_chars:
remaining = max_chars - total
if remaining > 200:
parts.append(chunk[:remaining] + "\n... [truncado por límite de tokens]")
break
parts.append(chunk)
total += len(chunk)
return "\n---\n".join(parts)
# ── Hierarchical (Map-Reduce) synthesis ─────────────────────────────
async def _map_chunk(
self,
chunk_docs: List[Dict[str, Any]],
chunk_idx: int,
query: str,
profile: str,
geo_context: str = "Automático",
) -> str:
"""Map step: synthesize a single chunk of documents."""
profile_data = AGENT_PROFILES.get(profile, AGENT_PROFILES["general"])
context = self.build_full_text_context(chunk_docs, max_chars=40000)
system = f"Eres un {profile_data['title']}. Sintetiza este bloque de documentos."
user = f"""CONSULTA ORIGINAL: "{query}"
DOCUMENTOS DEL BLOQUE {chunk_idx}:
{context}
TAREA: Sintetiza los hallazgos clave de este bloque.
- Menciona autores, años y datos específicos.
- Usa formato [[n]] {{BIB:ID}} para cada cita.
- Sé conciso pero técnico.
- SOLO texto, NO JSON."""
user += f"\n\nCONTEXTO GEOGRÁFICO A PRIORIZAR: {geo_context}"
return await self._call_llm(system, user, temperature=0.0)
async def _reduce_summaries(self, summaries: List[str], query: str, profile: str, geo_context: str = "Automático") -> dict:
"""Reduce step: merge chunk summaries into a single master plan."""
profile_data = AGENT_PROFILES.get(profile, AGENT_PROFILES["general"])
combined = "\n\n---\n\n".join(summaries)
system = f"Eres un {profile_data['title']}. Fusiona múltiples síntesis parciales en un plan coherente."
user = MASTER_SYNTHESIS_PROMPT.format(
query=query,
agent_title=profile_data["title"],
agent_title_upper=profile_data["title"].upper(),
profile_instruction=profile_data["instruction"],
template_structure="Integra las secciones de las síntesis parciales en un plan maestro unificado.",
)
user += f"\n\nSÍNTESIS PARCIALES:\n{combined}"
user += f"\n\nCONTEXTO GEOGRÁFICO A PRIORIZAR: {geo_context}"
response = await self._call_llm(system, user, temperature=0.0)
return extract_research_plan(response)
async def hierarchical_synthesis(
self,
query: str,
documents: List[Dict[str, Any]],
profile: str = "general",
chunk_size: int = 10,
geo_context: str = "Automático",
) -> dict:
"""
Map-Reduce hierarchical synthesis.
1. Split docs into chunks.
2. Map: synthesize each chunk independently.
3. Reduce: merge all chunk summaries into a master plan.
4. Detect gaps and optionally rescue.
"""
enriched = self.classify_evidence(documents)
sorted_docs = self.sort_by_evidence(enriched)
chunks = [
sorted_docs[i:i + chunk_size]
for i in range(0, len(sorted_docs), chunk_size)
]
summaries: List[str] = []
for idx, chunk in enumerate(chunks, 1):
summary = await self._map_chunk(chunk, idx, query, profile, geo_context=geo_context)
summaries.append(summary)
master_plan = await self._reduce_summaries(summaries, query, profile, geo_context=geo_context)
plan_sections = master_plan.get("plan", [])
gap_result = await self.detect_gaps(query, plan_sections)
master_plan["gap_analysis"] = gap_result
master_plan["evidence_summary"] = self.evidence_summary(enriched)
if gap_result.get("requires_rescue"):
rescue_result = await self._rescue_search(query, gap_result.get("missing_aspects", []))
master_plan["rescue_results"] = rescue_result
return master_plan
# ── Linear synthesis (original approach, kept for compatibility) ────
async def linear_synthesis(
self,
query: str,
documents: List[Dict[str, Any]],
profile: str = "general",
) -> dict:
"""Original linear pipeline: orchestrate → plan → master plan → gap detection."""
enriched = self.classify_evidence(documents)
sorted_docs = self.sort_by_evidence(enriched)
docs_context = self.build_full_text_context(sorted_docs)
master_plan = await self.generate_master_plan(query, docs_context, profile)
plan_sections = master_plan.get("plan", [])
gap_result = await self.detect_gaps(query, plan_sections)
master_plan["gap_analysis"] = gap_result
master_plan["evidence_summary"] = self.evidence_summary(enriched)
if gap_result.get("requires_rescue"):
rescue_result = await self._rescue_search(query, gap_result.get("missing_aspects", []))
master_plan["rescue_results"] = rescue_result
return master_plan
# ── Gap detection + rescue search ───────────────────────────────────
async def _rescue_search(self, query: str, missing_aspects: List[str]) -> Dict[str, Any]:
"""Generate supplementary search queries for detected gaps."""
system = "Eres un Estratega de Búsqueda de Rescate. Genera queries de búsqueda para cubrir faltas."
aspects_text = "\n".join(f"- {a}" for a in missing_aspects)
user = f"""CONSULTA ORIGINAL: "{query}"
ASPECTOS FALTANTES:
{aspects_text}
Genera queries de búsqueda de rescate optimizados para cubrir cada aspecto faltante.
RESPONDE EN JSON:
{{
"rescue_queries": [
{{"aspect": "...", "english_query": "...", "spanish_query": "..."}}
]
}}"""
response = await self._call_llm(system, user, temperature=0.0)
return self._parse_json(response)
async def run_full_pipeline(
self,
query: str,
documents: List[Dict[str, Any]],
profile: str = "general",
mode: str = "linear",
chunk_size: int = 10,
geo_context: str = "Automático",
) -> dict:
"""
Unified entry point for the full synthesis pipeline.
mode: "linear" | "hierarchical"
"""
if mode == "hierarchical":
return await self.hierarchical_synthesis(query, documents, profile, chunk_size, geo_context=geo_context)
return await self.linear_synthesis(query, documents, profile)
# ── Cleanup ─────────────────────────────────────────────────────────
async def close(self):
await self.client.aclose()