Módulo 7: Flujos Avanzados

Proyecto Evolutivo: Retry Logic y Branching (v2)

Descripción del proyecto

En el Módulo 6, construiste la v1 del AI Research Assistant: un agente que descompone un tema en sub-queries, busca en 3 fuentes, sintetiza hallazgos, y genera un reporte estructurado con Pydantic. Funciona end-to-end. Pero tiene un problema fundamental: es frágil.

La v1 busca en 3 fuentes para cada sub-query usando futures. Eso parece paralelo — y técnicamente lo es dentro de search_all_sources. Pero las búsquedas por sub-query se lanzan como futures independientes sin protección: si la API falla, el pipeline crashea. Si una fuente tarda 30 segundos, todo espera. Si dos fuentes retornan el mismo hallazgo, se duplica en el reporte. No hay reintentos, no hay degradación graceful, y no hay visibilidad de qué pasó cuando algo falla.

La v2 que construyes en este módulo resuelve todo eso. No es una reescritura — es una evolución. Tu código de la v1 sigue funcionando. Este módulo le agrega capas de robustez encima: retry con backoff exponencial para que errores transitorios se recuperen solos, ejecución paralela protegida para que las 3 fuentes se busquen simultáneamente pero cada una maneje sus propios errores, merge con deduplicación para consolidar resultados inteligentemente, y logging estructurado para que sepas exactamente qué pasó en cada paso.

La diferencia es tangible. Si ejecutas v1 con una API que falla intermitentemente, tu agente crashea. Si ejecutas v2 con la misma API, tu agente reintenta, recupera, y si la API sigue caída después de 3 intentos, genera el reporte con las fuentes que sí respondieron. La v1 es un demo. La v2 es un sistema que puedes poner en producción.


Objetivo del proyecto

Evolucionar el AI Research Assistant de v1 (funcional pero frágil) a v2 (robusto y production-ready), agregando retry logic, búsqueda paralela protegida, merge con deduplicación, graceful degradation, y logging estructurado.

Al completar este proyecto:

  • 🔧 Implementarás retry con backoff exponencial y jitter para cada búsqueda
  • 🔧 Protegerás cada fuente de búsqueda con manejo de errores individual
  • 🔧 Ejecutarás búsquedas en paralelo con futures y recolectarás resultados con tolerancia a fallos
  • 🔧 Implementarás merge inteligente con deduplicación de hallazgos
  • 🔧 Agregarás graceful degradation: mínimo 1 fuente para un reporte válido
  • 🔧 Integrarás logging estructurado para trazabilidad de cada operación

Antes y después

v1 (Módulo 6): funcional pero frágil

Topic → Decompose → [search_all_sources(q1), search_all_sources(q2), ...] → Synthesize → Report

Problemas:
- Si search_web falla → Exception, pipeline crashea
- Si search_academic tarda 30s → todo espera 30s
- Si web y news retornan el mismo hallazgo → duplicado en el reporte
- Si algo falla → "Error durante la investigación" (sin detalles)

v2 (Este módulo): robusto y production-ready

Topic → Decompose → [search_with_retry(web), search_with_retry(academic), search_with_retry(news)] (paralelo)
                      ↓ cada una con retry 3x, backoff 1s→2s→4s
                      ↓ si falla después de 3 retries → graceful skip
                    → Merge con deduplicación
                    → Synthesize (con fuentes disponibles)
                    → Report (incluye status de cada fuente)

Mejoras:
- Error transitorio → auto-recovery (retry con backoff)
- API lenta → timeout + retry, no bloquea otros
- Hallazgos duplicados → deduplicación por similitud
- Fuente caída permanente → continúa con las demás
- Cualquier fallo → log estructurado con request_id, nodo, duración

Especificaciones técnicas

Stack tecnológico

ComponenteVersiónPropósito
Python3.11+Runtime
LangChainv1.2+Framework de LLMs
LangGraphv1.0+Functional API (@entrypoint, @task)
langchain-openailatestProveedor de modelos
pydanticv2+Modelos de datos structured
python-dotenvlatestVariables de entorno

Estructura del proyecto (evolución desde v1)

research-assistant/
├── .env
├── requirements.txt
├── agents/
│   └── researcher.py           # MODIFICADO — @entrypoint v2 con retry + parallel
├── tools/
│   ├── web_search.py           # MODIFICADO — search con simulación de fallos
│   └── calculator.py           # SIN CAMBIOS
├── state/
│   └── research_state.py       # EXTENDIDO — nuevos campos para source status
├── config/
│   └── settings.py             # EXTENDIDO — config de retry y timeouts
├── utils/                      # NUEVO
│   ├── retry.py                # Retry con backoff exponencial
│   └── logger.py               # Logging estructurado
└── main.py                     # MODIFICADO — muestra source availability

Los archivos nuevos son utils/retry.py y utils/logger.py. El resto son extensiones de archivos existentes. Tu código de v1 no se borra — se mejora.


Paso 1: Retry con backoff exponencial (utils/retry.py)

El primer building block. Cuando una API falla con un error transitorio (timeout, 429, 503), no quieres reintentar inmediatamente — eso empeora el problema. Quieres esperar un tiempo creciente: 1s, 2s, 4s. Y agregas jitter (variación aleatoria) para evitar que 100 clientes reintenten al mismo tiempo (thundering herd).

"""
utils/retry.py
Retry con backoff exponencial y jitter para el AI Research Assistant v2.
"""

import time
import random


class RetryConfig:
    """Configuración de retry con backoff exponencial."""

    def __init__(
        self,
        max_retries: int = 3,
        base_delay: float = 1.0,
        max_delay: float = 10.0,
        jitter: bool = True,
    ):
        self.max_retries = max_retries
        self.base_delay = base_delay
        self.max_delay = max_delay
        self.jitter = jitter

    def get_delay(self, attempt: int) -> float:
        """Calcula el delay para un intento específico."""
        delay = min(self.base_delay * (2 ** attempt), self.max_delay)
        if self.jitter:
            delay = delay * (0.5 + random.random() * 0.5)
        return delay


def retry_with_backoff(
    func,
    args: tuple = (),
    kwargs: dict = None,
    config: RetryConfig = None,
    on_retry=None,
    retriable_exceptions: tuple = (Exception,),
) -> dict:
    """
    Ejecuta func con retry y backoff exponencial.

    Retorna un dict con el resultado o el error final:
    - {"status": "success", "result": ..., "attempts": N}
    - {"status": "failed", "error": "...", "attempts": N}
    """
    if kwargs is None:
        kwargs = {}
    if config is None:
        config = RetryConfig()

    last_error = None

    for attempt in range(config.max_retries + 1):
        try:
            result = func(*args, **kwargs)
            return {
                "status": "success",
                "result": result,
                "attempts": attempt + 1,
            }
        except retriable_exceptions as e:
            last_error = e

            if attempt < config.max_retries:
                delay = config.get_delay(attempt)
                if on_retry:
                    on_retry(attempt + 1, config.max_retries, delay, str(e))
                time.sleep(delay)
            else:
                return {
                    "status": "failed",
                    "error": str(e),
                    "attempts": attempt + 1,
                }

    return {
        "status": "failed",
        "error": str(last_error),
        "attempts": config.max_retries + 1,
    }

El patrón de backoff: intento 0 = sin espera, intento 1 = ~1s, intento 2 = ~2s, intento 3 = ~4s. El jitter hace que el delay real varíe entre 50% y 100% del valor calculado. Así, 100 clientes que fallan al mismo tiempo no reintentan al mismo tiempo.

Verifiquemos que funciona:

from utils.retry import retry_with_backoff, RetryConfig

call_count = 0

def flaky_function():
    """Falla las primeras 2 veces, funciona a la tercera."""
    global call_count
    call_count += 1
    if call_count <= 2:
        raise ConnectionError(f"Intento {call_count}: API no disponible")
    return "¡Éxito en el intento 3!"


config = RetryConfig(max_retries=3, base_delay=0.5, jitter=False)
result = retry_with_backoff(
    flaky_function,
    config=config,
    on_retry=lambda attempt, max_r, delay, err: print(
        f"  Retry {attempt}/{max_r} en {delay:.1f}s — {err}"
    ),
)
print(f"Status: {result['status']}")
print(f"Resultado: {result.get('result', result.get('error'))}")
print(f"Intentos: {result['attempts']}")
# Output esperado:
#   Retry 1/3 en 0.5s — Intento 1: API no disponible
#   Retry 2/3 en 1.0s — Intento 2: API no disponible
# Status: success
# Resultado: ¡Éxito en el intento 3!
# Intentos: 3

Paso 2: Logging estructurado (utils/logger.py)

Cada operación del agente necesita ser rastreable. El logger emite JSON con request_id, nombre del nodo, duración, status, y cualquier metadata relevante.

"""
utils/logger.py
Logging estructurado para el AI Research Assistant v2.
"""

import json
import time
import logging

logging.basicConfig(
    level=logging.INFO,
    format="%(message)s",
)


class ResearchLogger:
    """Logger estructurado con correlation ID para el Research Assistant."""

    def __init__(self, name: str = "research_agent"):
        self._logger = logging.getLogger(name)
        self._request_id = None

    def set_request_id(self, request_id: str):
        self._request_id = request_id

    def _emit(self, level: str, node: str, event: str, **kwargs):
        entry = {
            "ts": time.strftime("%Y-%m-%dT%H:%M:%S"),
            "level": level,
            "request_id": self._request_id,
            "node": node,
            "event": event,
        }
        entry.update(kwargs)
        self._logger.info(json.dumps(entry, ensure_ascii=False))

    def start(self, node: str, **kwargs):
        self._emit("INFO", node, "start", **kwargs)

    def end(self, node: str, duration_ms: float, **kwargs):
        self._emit("INFO", node, "end", duration_ms=round(duration_ms, 1), **kwargs)

    def retry(self, node: str, attempt: int, max_retries: int, delay: float, error: str, **kwargs):
        self._emit("WARN", node, "retry", attempt=attempt, max_retries=max_retries, delay_s=round(delay, 2), error=error, **kwargs)

    def error(self, node: str, error: str, duration_ms: float, **kwargs):
        self._emit("ERROR", node, "error", error=error, duration_ms=round(duration_ms, 1), **kwargs)

    def skip(self, node: str, reason: str, **kwargs):
        self._emit("WARN", node, "skip", reason=reason, **kwargs)

Paso 3: Extender los modelos de datos (state/research_state.py)

El reporte v2 necesita información sobre qué fuentes respondieron y cuáles fallaron. Extiende los modelos Pydantic — sin romper la v1:

"""
state/research_state.py
Modelos de datos para el AI Research Assistant v2.
Extiende v1 con campos de source availability.
"""

from pydantic import BaseModel, Field
from datetime import datetime


class Source(BaseModel):
    """Una fuente de información encontrada durante la investigación."""
    name: str = Field(description="Nombre de la fuente")
    source_type: str = Field(description="Tipo: web, academic, news")
    content: str = Field(description="Contenido extraído de la fuente")


class KeyFinding(BaseModel):
    """Un hallazgo clave identificado en la investigación."""
    title: str = Field(description="Título del hallazgo")
    description: str = Field(description="Descripción del hallazgo")
    confidence: float = Field(
        description="Confianza en el hallazgo (0.0 a 1.0)",
        ge=0.0,
        le=1.0,
    )


class SourceStatus(BaseModel):
    """Estado de una fuente de búsqueda (v2)."""
    source_type: str = Field(description="Tipo de fuente")
    status: str = Field(description="ok, failed, timeout, skipped")
    attempts: int = Field(default=1, description="Número de intentos realizados")
    error: str = Field(default="", description="Mensaje de error si falló")
    duration_ms: float = Field(default=0, description="Duración total incluyendo retries")


class ResearchReport(BaseModel):
    """Reporte estructurado de investigación (v2)."""
    topic: str = Field(description="Tema investigado")
    summary: str = Field(description="Resumen ejecutivo (2-3 oraciones)")
    key_findings: list[KeyFinding] = Field(
        description="Hallazgos principales",
        min_length=1,
    )
    sources: list[Source] = Field(
        description="Fuentes consultadas exitosamente",
        min_length=1,
    )
    sub_queries: list[str] = Field(
        description="Sub-preguntas generadas para la investigación",
    )
    confidence: float = Field(
        description="Confianza general del reporte (0.0 a 1.0)",
        ge=0.0,
        le=1.0,
    )
    generated_at: str = Field(
        default_factory=lambda: datetime.now().isoformat(),
        description="Timestamp de generación",
    )
    source_availability: list[SourceStatus] = Field(
        default_factory=list,
        description="(v2) Estado de cada fuente de búsqueda",
    )
    version: str = Field(default="v2", description="Versión del agente")


class SubQuery(BaseModel):
    """Una sub-pregunta generada a partir del tema principal."""
    query: str = Field(description="La sub-pregunta")
    rationale: str = Field(description="Por qué esta pregunta es relevante")

Los cambios:

  • SourceStatus — modelo nuevo para trackear qué pasó con cada fuente
  • source_availability en ResearchReport — lista de SourceStatus con default vacío (compatible con v1)
  • version — distingue reportes v1 de v2

Todos los cambios son aditivos. Un reporte v1 sigue siendo un ResearchReport válido — los campos nuevos tienen defaults.


Paso 4: Extender la configuración (config/settings.py)

Agrega configuración de retry y timeouts:

"""
config/settings.py
Configuración del AI Research Assistant v2.
"""

from dotenv import load_dotenv
load_dotenv()

MODEL_NAME = "openai:gpt-4.1-mini"
MODEL_TEMPERATURE = 0.2

MAX_SUB_QUERIES = 4
MAX_SOURCES_PER_QUERY = 3

SEARCH_SOURCES = ["web", "academic", "news"]

# v2: Retry configuration
RETRY_MAX_ATTEMPTS = 3
RETRY_BASE_DELAY = 1.0
RETRY_MAX_DELAY = 10.0
RETRY_JITTER = True

# v2: Minimum sources for a valid report
MIN_SOURCES_FOR_REPORT = 1

# v2: Search timeout
SEARCH_TIMEOUT_SECONDS = 5.0

Paso 5: Evolucionar las tools de búsqueda (tools/web_search.py)

Las tools de la v1 nunca fallaban — eran mocks determinísticos. La v2 simula errores reales para que el retry logic tenga algo contra lo cual trabajar. Pero el mock sigue siendo controlable para testing.

"""
tools/web_search.py
Mock web search tool para el AI Research Assistant v2.
Simula errores transitorios y latencia variable para testing de retry logic.
"""

import time
import random
import hashlib
from langgraph.func import task

from utils.retry import retry_with_backoff, RetryConfig
from utils.logger import ResearchLogger


MOCK_RESULTS = {
    "web": {
        "default": (
            "Múltiples fuentes web coinciden en que {query} es un tema de creciente "
            "interés. Expertos destacan avances significativos en los últimos 2 años. "
            "Las aplicaciones prácticas incluyen automatización, análisis de datos "
            "y generación de contenido."
        ),
    },
    "academic": {
        "default": (
            "Investigación reciente (2025-2026) muestra que {query} tiene fundamentos "
            "teóricos sólidos respaldados por múltiples estudios peer-reviewed. "
            "Los resultados experimentales demuestran mejoras del 40-60% en métricas "
            "clave comparados con métodos tradicionales."
        ),
    },
    "news": {
        "default": (
            "Noticias recientes reportan que {query} está generando impacto en la "
            "industria. Empresas líderes como Google, Microsoft y startups emergentes "
            "están invirtiendo significativamente en esta área. Se esperan desarrollos "
            "importantes para finales de 2026."
        ),
    },
}

FAILURE_CONFIG = {
    "web": {"fail_rate": 0.0, "latency_range": (0.1, 0.3)},
    "academic": {"fail_rate": 0.0, "latency_range": (0.2, 0.5)},
    "news": {"fail_rate": 0.0, "latency_range": (0.1, 0.4)},
}


def configure_failures(source: str, fail_rate: float = 0.0, latency_range: tuple = None):
    """Configura la tasa de fallos de una fuente para testing."""
    if latency_range is None:
        latency_range = (0.1, 0.3)
    FAILURE_CONFIG[source] = {
        "fail_rate": fail_rate,
        "latency_range": latency_range,
    }


def _generate_deterministic_score(query: str, source_type: str) -> float:
    hash_input = f"{query}:{source_type}"
    hash_value = int(hashlib.md5(hash_input.encode()).hexdigest()[:8], 16)
    return round(0.5 + (hash_value % 50) / 100, 2)


def _raw_search(query: str, source_type: str) -> dict:
    """Búsqueda raw que puede fallar según FAILURE_CONFIG."""
    config = FAILURE_CONFIG.get(source_type, {"fail_rate": 0.0, "latency_range": (0.1, 0.3)})

    latency = random.uniform(*config["latency_range"])
    time.sleep(latency)

    if random.random() < config["fail_rate"]:
        error_types = [
            ConnectionError(f"{source_type} API: connection refused"),
            TimeoutError(f"{source_type} API: request timed out"),
            RuntimeError(f"{source_type} API: 503 Service Unavailable"),
        ]
        raise random.choice(error_types)

    content = MOCK_RESULTS[source_type]["default"].format(query=query)
    return {
        "source_name": f"{source_type.title()} Search",
        "source_type": source_type,
        "content": content,
        "relevance": _generate_deterministic_score(query, source_type),
    }


@task
def search_with_retry(query: str, source_type: str, logger: ResearchLogger = None) -> dict:
    """
    Búsqueda con retry y backoff exponencial (v2).
    Retorna resultado exitoso o error con metadata de intentos.
    """
    start = time.time()
    retry_config = RetryConfig(
        max_retries=3,
        base_delay=1.0,
        max_delay=10.0,
        jitter=True,
    )

    def on_retry(attempt, max_retries, delay, error):
        if logger:
            logger.retry(
                f"search_{source_type}",
                attempt=attempt,
                max_retries=max_retries,
                delay=delay,
                error=error,
                query=query[:50],
            )

    outcome = retry_with_backoff(
        func=_raw_search,
        args=(query, source_type),
        config=retry_config,
        on_retry=on_retry,
        retriable_exceptions=(ConnectionError, TimeoutError, RuntimeError),
    )

    duration_ms = (time.time() - start) * 1000

    if outcome["status"] == "success":
        if logger:
            logger.end(
                f"search_{source_type}",
                duration_ms=duration_ms,
                attempts=outcome["attempts"],
                status="success",
                query=query[:50],
            )
        return {
            **outcome["result"],
            "search_status": "ok",
            "attempts": outcome["attempts"],
            "duration_ms": round(duration_ms, 1),
        }
    else:
        if logger:
            logger.error(
                f"search_{source_type}",
                error=outcome["error"],
                duration_ms=duration_ms,
                attempts=outcome["attempts"],
                query=query[:50],
            )
        return {
            "source_name": f"{source_type.title()} Search",
            "source_type": source_type,
            "content": "",
            "relevance": 0.0,
            "search_status": "failed",
            "error": outcome["error"],
            "attempts": outcome["attempts"],
            "duration_ms": round(duration_ms, 1),
        }


SEARCH_FUNCTIONS = {
    "web": search_with_retry,
    "academic": search_with_retry,
    "news": search_with_retry,
}

Cambios respecto a v1:

  • _raw_search simula errores transitorios configurables con FAILURE_CONFIG
  • search_with_retry envuelve _raw_search con retry + backoff + logging
  • ✅ Nunca lanza excepciones — retorna {"search_status": "ok"} o {"search_status": "failed"}
  • configure_failures() permite activar/desactivar fallos para testing

Paso 6: Evolucionar el agente principal (agents/researcher.py)

Este es el cambio central. El @entrypoint ahora ejecuta búsquedas con retry, merge con deduplicación, y graceful degradation.

"""
agents/researcher.py
Agente principal del AI Research Assistant (v2).
Evolución de v1: retry, parallel con protección, merge, graceful degradation.
"""

import json
import time
import uuid
from langchain.chat_models import init_chat_model
from langgraph.func import entrypoint, task
from langgraph.checkpoint.memory import MemorySaver

import sys
sys.path.insert(0, ".")

from config.settings import (
    MODEL_NAME,
    MODEL_TEMPERATURE,
    MAX_SUB_QUERIES,
    SEARCH_SOURCES,
    MIN_SOURCES_FOR_REPORT,
)
from state.research_state import (
    ResearchReport,
    Source,
    KeyFinding,
    SubQuery,
    SourceStatus,
)
from tools.web_search import search_with_retry
from tools.calculator import calculate_confidence
from utils.logger import ResearchLogger


model = init_chat_model(MODEL_NAME, temperature=MODEL_TEMPERATURE)
agent_logger = ResearchLogger("research_agent_v2")


# =============================================================================
# TASK: Descomponer tema en sub-queries (sin cambios de v1)
# =============================================================================

@task
def decompose_query(topic: str) -> list[dict]:
    """Descompone un tema de investigación en sub-preguntas específicas."""
    response = model.invoke(
        f"Eres un investigador experto. Descompone este tema en "
        f"{MAX_SUB_QUERIES} sub-preguntas específicas e investigables.\n\n"
        f"Tema: {topic}\n\n"
        f"Responde en JSON (sin markdown, sin ```json):\n"
        f'[{{"query": "sub-pregunta", "rationale": "por qué es relevante"}}]\n\n'
        f"Solo el JSON, nada más."
    )

    try:
        queries = json.loads(response.content)
        return queries[:MAX_SUB_QUERIES]
    except json.JSONDecodeError:
        return [
            {"query": topic, "rationale": "Query original como fallback"},
            {"query": f"avances recientes en {topic}", "rationale": "Tendencias actuales"},
            {"query": f"aplicaciones prácticas de {topic}", "rationale": "Uso real"},
        ]


# =============================================================================
# TASK: Buscar en todas las fuentes con retry (v2)
# =============================================================================

@task
def search_all_sources_v2(query: str, logger: ResearchLogger) -> list[dict]:
    """
    Busca en todas las fuentes con retry individual por fuente (v2).
    Lanza todas las búsquedas en paralelo con futures.
    Cada búsqueda tiene su propio retry — un fallo en 'academic' no afecta a 'web'.
    """
    logger.start("search_all", query=query[:50], sources=SEARCH_SOURCES)
    start = time.time()

    futures = []
    for source_type in SEARCH_SOURCES:
        logger.start(f"search_{source_type}", query=query[:50])
        futures.append(search_with_retry(query, source_type, logger))

    results = [f.result() for f in futures]

    duration_ms = (time.time() - start) * 1000
    ok_count = sum(1 for r in results if r["search_status"] == "ok")
    logger.end("search_all", duration_ms=duration_ms, sources_ok=ok_count, sources_total=len(results))

    return results


# =============================================================================
# TASK: Merge con deduplicación (v2 — NUEVO)
# =============================================================================

@task
def merge_and_deduplicate(all_results: list[dict]) -> list[dict]:
    """
    Combina resultados de múltiples búsquedas y elimina duplicados.
    Deduplicación por: mismo source_type + contenido muy similar.
    """
    seen_keys = set()
    unique_results = []

    for result in all_results:
        if result["search_status"] != "ok":
            continue

        dedup_key = f"{result['source_type']}:{result['content'][:100]}"

        if dedup_key not in seen_keys:
            seen_keys.add(dedup_key)
            unique_results.append(result)

    unique_results.sort(key=lambda r: r.get("relevance", 0), reverse=True)

    return unique_results


# =============================================================================
# TASK: Sintetizar hallazgos (sin cambios de v1)
# =============================================================================

@task
def synthesize_findings(topic: str, all_results: list[dict]) -> list[dict]:
    """Sintetiza resultados de búsqueda en hallazgos clave."""
    results_text = ""
    for i, result in enumerate(all_results, 1):
        results_text += (
            f"\nFuente {i} ({result['source_type']}): {result['content']}\n"
        )

    response = model.invoke(
        f"Eres un analista de investigación. Basándote en estas fuentes, "
        f"identifica 3-5 hallazgos clave sobre '{topic}'.\n\n"
        f"Fuentes:\n{results_text}\n\n"
        f"Responde en JSON (sin markdown, sin ```json):\n"
        f'[{{"title": "título corto", "description": "descripción de 1-2 oraciones", '
        f'"confidence": 0.8}}]\n\n'
        f"confidence es de 0.0 a 1.0. Solo el JSON, nada más."
    )

    try:
        findings = json.loads(response.content)
        return findings[:5]
    except json.JSONDecodeError:
        return [{
            "title": "Hallazgo general",
            "description": f"La investigación sobre {topic} muestra resultados relevantes en múltiples fuentes.",
            "confidence": 0.6,
        }]


# =============================================================================
# TASK: Generar resumen ejecutivo (sin cambios de v1)
# =============================================================================

@task
def generate_summary(topic: str, findings: list[dict], source_info: str) -> str:
    """Genera un resumen ejecutivo que incluye info sobre disponibilidad de fuentes."""
    findings_text = "\n".join(
        f"- {f['title']}: {f['description']}" for f in findings
    )

    response = model.invoke(
        f"Genera un resumen ejecutivo de 2-3 oraciones sobre la investigación "
        f"del tema '{topic}'.\n\n"
        f"Hallazgos principales:\n{findings_text}\n\n"
        f"Nota sobre fuentes: {source_info}\n\n"
        f"Solo el resumen, sin título ni formato extra."
    )
    return response.content.strip()


# =============================================================================
# ENTRYPOINT: Agente de investigación v2
# =============================================================================

memory = MemorySaver()


@entrypoint(checkpointer=memory)
def research_agent(topic: str) -> dict:
    """
    AI Research Assistant v2.
    Flujo: topic → decompose → search (parallel + retry) → merge → synthesize → report.
    """
    request_id = uuid.uuid4().hex[:8]
    agent_logger.set_request_id(request_id)
    agent_logger.start("pipeline", topic=topic, version="v2")

    pipeline_start = time.time()

    print(f"\n{'=' * 60}")
    print(f"  🔬 AI Research Assistant v2")
    print(f"  Tema: {topic}")
    print(f"  Request ID: {request_id}")
    print(f"{'=' * 60}")

    # --- Paso 1: Descomponer el tema ---
    print(f"\n📋 Paso 1: Descomponiendo tema en sub-queries...")
    agent_logger.start("decompose", topic=topic[:50])
    decompose_start = time.time()

    sub_queries_raw = decompose_query(topic).result()
    sub_queries = [SubQuery(**sq) for sq in sub_queries_raw]

    agent_logger.end("decompose", duration_ms=(time.time() - decompose_start) * 1000, count=len(sub_queries))
    print(f"   ✓ {len(sub_queries)} sub-queries generadas:")
    for i, sq in enumerate(sub_queries, 1):
        print(f"     {i}. {sq.query}")

    # --- Paso 2: Buscar en paralelo con retry ---
    print(f"\n🔍 Paso 2: Buscando en {len(SEARCH_SOURCES)} fuentes por sub-query (con retry)...")

    all_raw_results = []
    source_statuses = []

    search_futures = [
        search_all_sources_v2(sq.query, agent_logger) for sq in sub_queries
    ]

    for i, future in enumerate(search_futures):
        results = future.result()
        all_raw_results.extend(results)

        ok = sum(1 for r in results if r["search_status"] == "ok")
        failed = sum(1 for r in results if r["search_status"] == "failed")
        print(f"   Sub-query {i + 1}: {ok} OK, {failed} failed")

        for r in results:
            source_statuses.append(SourceStatus(
                source_type=r["source_type"],
                status=r["search_status"],
                attempts=r.get("attempts", 1),
                error=r.get("error", ""),
                duration_ms=r.get("duration_ms", 0),
            ))

    total_ok = sum(1 for r in all_raw_results if r["search_status"] == "ok")
    total_failed = sum(1 for r in all_raw_results if r["search_status"] == "failed")
    print(f"   Total: {total_ok} exitosas, {total_failed} fallidas de {len(all_raw_results)}")

    # --- Verificar mínimo de fuentes ---
    if total_ok < MIN_SOURCES_FOR_REPORT:
        pipeline_ms = (time.time() - pipeline_start) * 1000
        agent_logger.error("pipeline", error="Insufficient sources", duration_ms=pipeline_ms, sources_ok=total_ok)
        print(f"\n   ❌ Fuentes insuficientes ({total_ok} < {MIN_SOURCES_FOR_REPORT}). Abortando.")

        return {
            "error": f"Solo {total_ok} fuentes respondieron. Mínimo requerido: {MIN_SOURCES_FOR_REPORT}",
            "request_id": request_id,
            "source_availability": [s.model_dump() for s in source_statuses],
            "version": "v2",
        }

    # --- Paso 3: Merge con deduplicación ---
    print(f"\n🔀 Paso 3: Merge y deduplicación de resultados...")
    unique_results = merge_and_deduplicate(all_raw_results).result()
    print(f"   ✓ {len(all_raw_results)} resultados raw → {len(unique_results)} únicos después de deduplicación")

    # --- Paso 4: Sintetizar hallazgos ---
    print(f"\n🧠 Paso 4: Sintetizando hallazgos...")
    findings_raw = synthesize_findings(topic, unique_results).result()
    findings = [KeyFinding(**f) for f in findings_raw]
    print(f"   ✓ {len(findings)} hallazgos identificados:")
    for i, f in enumerate(findings, 1):
        print(f"     {i}. [{f.confidence:.0%}] {f.title}")

    # --- Paso 5: Generar resumen ---
    print(f"\n📝 Paso 5: Generando resumen ejecutivo...")
    source_info = f"{total_ok} de {total_ok + total_failed} fuentes respondieron exitosamente"
    if total_failed > 0:
        failed_sources = set(r["source_type"] for r in all_raw_results if r["search_status"] == "failed")
        source_info += f". Fuentes no disponibles: {', '.join(failed_sources)}"

    summary = generate_summary(topic, findings_raw, source_info).result()
    print(f"   ✓ Resumen generado ({len(summary)} chars)")

    # --- Paso 6: Calcular confianza (ajustada por availability) ---
    print(f"\n📊 Paso 6: Calculando confianza del reporte...")
    avg_relevance = (
        sum(r["relevance"] for r in unique_results) / len(unique_results)
        if unique_results else 0.5
    )

    base_confidence = calculate_confidence(
        num_sources=len(unique_results),
        avg_relevance=avg_relevance,
        num_findings=len(findings),
    ).result()

    availability_factor = total_ok / (total_ok + total_failed) if (total_ok + total_failed) > 0 else 0.5
    adjusted_confidence = round(base_confidence * (0.7 + 0.3 * availability_factor), 2)
    print(f"   ✓ Confianza base: {base_confidence:.0%}, ajustada por availability: {adjusted_confidence:.0%}")

    # --- Paso 7: Construir reporte structured ---
    print(f"\n📄 Paso 7: Construyendo reporte v2...")
    sources = [
        Source(
            name=r["source_name"],
            source_type=r["source_type"],
            content=r["content"],
        )
        for r in unique_results
    ]

    report = ResearchReport(
        topic=topic,
        summary=summary,
        key_findings=findings,
        sources=sources,
        sub_queries=[sq.query for sq in sub_queries],
        confidence=adjusted_confidence,
        source_availability=source_statuses,
        version="v2",
    )

    pipeline_ms = (time.time() - pipeline_start) * 1000
    agent_logger.end("pipeline", duration_ms=pipeline_ms, sources_ok=total_ok, sources_failed=total_failed, confidence=adjusted_confidence)

    print(f"   ✓ Reporte v2 generado exitosamente")
    print(f"   📊 Pipeline total: {pipeline_ms:.0f}ms")
    print(f"\n{'=' * 60}")

    return report.model_dump()

Las diferencias clave respecto a v1:

  • Request ID — cada ejecución tiene un ID único para trazabilidad
  • search_all_sources_v2 — cada fuente tiene retry individual, los errores no propagan
  • merge_and_deduplicate — paso nuevo que elimina duplicados y rankea por relevancia
  • Graceful degradation — verifica MIN_SOURCES_FOR_REPORT antes de sintetizar
  • Confianza ajustada — la confianza baja cuando hay fuentes no disponibles
  • source_availability — el reporte incluye el status de cada fuente
  • Logging estructurado — cada paso se logea con duración y metadata

Paso 7: Actualizar el CLI (main.py)

El CLI ahora muestra información de disponibilidad de fuentes y el request ID:

"""
main.py
CLI entrypoint para el AI Research Assistant v2.
"""

import sys
import json
import uuid

sys.path.insert(0, ".")

from agents.researcher import research_agent


def format_report(report: dict) -> str:
    """Formatea el reporte v2 para display en terminal."""
    if "error" in report:
        lines = [
            "",
            "╔" + "═" * 58 + "╗",
            "║" + "  ❌ ERROR EN INVESTIGACIÓN".center(58) + "║",
            "╚" + "═" * 58 + "╝",
            f"\n  {report['error']}",
            f"  Request ID: {report.get('request_id', 'N/A')}",
        ]
        if "source_availability" in report:
            lines.append(f"\n  Disponibilidad de fuentes:")
            for sa in report["source_availability"]:
                status_icon = "✅" if sa["status"] == "ok" else "❌"
                lines.append(f"    {status_icon} {sa['source_type']}: {sa['status']} ({sa['attempts']} intentos)")
        return "\n".join(lines)

    lines = []
    lines.append("")
    lines.append("╔" + "═" * 58 + "╗")
    lines.append("║" + f"  📄 REPORTE DE INVESTIGACIÓN ({report.get('version', 'v1')})".center(58) + "║")
    lines.append("╚" + "═" * 58 + "╝")

    lines.append(f"\n📌 Tema: {report['topic']}")
    lines.append(f"📅 Generado: {report['generated_at']}")
    lines.append(f"🎯 Confianza: {report['confidence']:.0%}")

    lines.append(f"\n{'─' * 60}")
    lines.append("📋 RESUMEN EJECUTIVO")
    lines.append(f"{'─' * 60}")
    lines.append(report["summary"])

    lines.append(f"\n{'─' * 60}")
    lines.append("🔍 SUB-QUERIES INVESTIGADAS")
    lines.append(f"{'─' * 60}")
    for i, sq in enumerate(report["sub_queries"], 1):
        lines.append(f"  {i}. {sq}")

    lines.append(f"\n{'─' * 60}")
    lines.append("💡 HALLAZGOS PRINCIPALES")
    lines.append(f"{'─' * 60}")
    for i, finding in enumerate(report["key_findings"], 1):
        conf = finding["confidence"]
        lines.append(f"\n  {i}. {finding['title']} [{conf:.0%} confianza]")
        lines.append(f"     {finding['description']}")

    lines.append(f"\n{'─' * 60}")
    lines.append(f"📚 FUENTES CONSULTADAS ({len(report['sources'])})")
    lines.append(f"{'─' * 60}")
    seen = set()
    for source in report["sources"]:
        key = f"{source['name']}:{source['source_type']}"
        if key not in seen:
            seen.add(key)
            lines.append(f"  • [{source['source_type'].upper()}] {source['name']}")

    if report.get("source_availability"):
        lines.append(f"\n{'─' * 60}")
        lines.append("🔌 DISPONIBILIDAD DE FUENTES")
        lines.append(f"{'─' * 60}")
        for sa in report["source_availability"]:
            icon = "✅" if sa["status"] == "ok" else "❌"
            retry_info = f" ({sa['attempts']} intentos, {sa['duration_ms']:.0f}ms)" if sa["attempts"] > 1 else f" ({sa['duration_ms']:.0f}ms)"
            error_info = f" — {sa['error']}" if sa.get("error") else ""
            lines.append(f"  {icon} {sa['source_type']}: {sa['status']}{retry_info}{error_info}")

    lines.append(f"\n{'═' * 60}")

    return "\n".join(lines)


def run_interactive():
    """Modo interactivo: el usuario escribe temas."""
    print("=" * 60)
    print("  🔬 AI Research Assistant v2")
    print("  Escribe un tema para investigar.")
    print("  Comandos: 'salir' para terminar")
    print("=" * 60)

    while True:
        try:
            topic = input("\n🔎 Tema: ").strip()
        except (KeyboardInterrupt, EOFError):
            print("\n\n¡Hasta luego!")
            break

        if not topic:
            continue

        if topic.lower() in ("salir", "exit", "quit"):
            print("\n¡Hasta luego!")
            break

        thread_id = f"research-v2-{uuid.uuid4().hex[:8]}"

        try:
            report = research_agent.invoke(
                topic,
                config={"configurable": {"thread_id": thread_id}},
            )
            print(format_report(report))

        except Exception as e:
            print(f"\n❌ Error inesperado: {e}")
            print("   Intenta con otro tema.")


def run_single(topic: str):
    """Ejecuta una sola investigación."""
    thread_id = f"research-v2-{uuid.uuid4().hex[:8]}"

    report = research_agent.invoke(
        topic,
        config={"configurable": {"thread_id": thread_id}},
    )
    print(format_report(report))

    print("\n📦 Reporte JSON:")
    print(json.dumps(report, indent=2, ensure_ascii=False))


if __name__ == "__main__":
    if len(sys.argv) > 1:
        run_single(" ".join(sys.argv[1:]))
    else:
        run_interactive()

Código completo actualizado

Para referencia rápida, aquí están todos los archivos del proyecto v2 consolidados.

utils/retry.py

"""
utils/retry.py
Retry con backoff exponencial y jitter para el AI Research Assistant v2.
"""

import time
import random


class RetryConfig:
    def __init__(
        self,
        max_retries: int = 3,
        base_delay: float = 1.0,
        max_delay: float = 10.0,
        jitter: bool = True,
    ):
        self.max_retries = max_retries
        self.base_delay = base_delay
        self.max_delay = max_delay
        self.jitter = jitter

    def get_delay(self, attempt: int) -> float:
        delay = min(self.base_delay * (2 ** attempt), self.max_delay)
        if self.jitter:
            delay = delay * (0.5 + random.random() * 0.5)
        return delay


def retry_with_backoff(
    func,
    args: tuple = (),
    kwargs: dict = None,
    config: RetryConfig = None,
    on_retry=None,
    retriable_exceptions: tuple = (Exception,),
) -> dict:
    if kwargs is None:
        kwargs = {}
    if config is None:
        config = RetryConfig()

    last_error = None

    for attempt in range(config.max_retries + 1):
        try:
            result = func(*args, **kwargs)
            return {"status": "success", "result": result, "attempts": attempt + 1}
        except retriable_exceptions as e:
            last_error = e
            if attempt < config.max_retries:
                delay = config.get_delay(attempt)
                if on_retry:
                    on_retry(attempt + 1, config.max_retries, delay, str(e))
                time.sleep(delay)

    return {"status": "failed", "error": str(last_error), "attempts": config.max_retries + 1}

utils/logger.py

"""
utils/logger.py
Logging estructurado para el AI Research Assistant v2.
"""

import json
import time
import logging

logging.basicConfig(level=logging.INFO, format="%(message)s")


class ResearchLogger:
    def __init__(self, name: str = "research_agent"):
        self._logger = logging.getLogger(name)
        self._request_id = None

    def set_request_id(self, request_id: str):
        self._request_id = request_id

    def _emit(self, level: str, node: str, event: str, **kwargs):
        entry = {
            "ts": time.strftime("%Y-%m-%dT%H:%M:%S"),
            "level": level,
            "request_id": self._request_id,
            "node": node,
            "event": event,
        }
        entry.update(kwargs)
        self._logger.info(json.dumps(entry, ensure_ascii=False))

    def start(self, node, **kw): self._emit("INFO", node, "start", **kw)
    def end(self, node, duration_ms, **kw): self._emit("INFO", node, "end", duration_ms=round(duration_ms, 1), **kw)
    def retry(self, node, attempt, max_retries, delay, error, **kw): self._emit("WARN", node, "retry", attempt=attempt, max_retries=max_retries, delay_s=round(delay, 2), error=error, **kw)
    def error(self, node, error, duration_ms, **kw): self._emit("ERROR", node, "error", error=error, duration_ms=round(duration_ms, 1), **kw)
    def skip(self, node, reason, **kw): self._emit("WARN", node, "skip", reason=reason, **kw)

state/research_state.py

"""
state/research_state.py
Modelos de datos para el AI Research Assistant v2.
"""

from pydantic import BaseModel, Field
from datetime import datetime


class Source(BaseModel):
    name: str = Field(description="Nombre de la fuente")
    source_type: str = Field(description="Tipo: web, academic, news")
    content: str = Field(description="Contenido extraído de la fuente")


class KeyFinding(BaseModel):
    title: str = Field(description="Título del hallazgo")
    description: str = Field(description="Descripción del hallazgo")
    confidence: float = Field(ge=0.0, le=1.0)


class SourceStatus(BaseModel):
    source_type: str = Field(description="Tipo de fuente")
    status: str = Field(description="ok, failed, timeout, skipped")
    attempts: int = Field(default=1)
    error: str = Field(default="")
    duration_ms: float = Field(default=0)


class ResearchReport(BaseModel):
    topic: str
    summary: str
    key_findings: list[KeyFinding] = Field(min_length=1)
    sources: list[Source] = Field(min_length=1)
    sub_queries: list[str]
    confidence: float = Field(ge=0.0, le=1.0)
    generated_at: str = Field(default_factory=lambda: datetime.now().isoformat())
    source_availability: list[SourceStatus] = Field(default_factory=list)
    version: str = Field(default="v2")


class SubQuery(BaseModel):
    query: str
    rationale: str

config/settings.py

"""
config/settings.py
Configuración del AI Research Assistant v2.
"""

from dotenv import load_dotenv
load_dotenv()

MODEL_NAME = "openai:gpt-4.1-mini"
MODEL_TEMPERATURE = 0.2

MAX_SUB_QUERIES = 4
MAX_SOURCES_PER_QUERY = 3
SEARCH_SOURCES = ["web", "academic", "news"]

RETRY_MAX_ATTEMPTS = 3
RETRY_BASE_DELAY = 1.0
RETRY_MAX_DELAY = 10.0
RETRY_JITTER = True

MIN_SOURCES_FOR_REPORT = 1
SEARCH_TIMEOUT_SECONDS = 5.0

Ejecución

Modo normal (sin fallos simulados)

cd research-assistant
python main.py "impacto de la inteligencia artificial en la educación"

Output esperado: todas las fuentes responden, 0 retries, confianza alta.

Modo con fallos simulados

Para probar retry logic, agrega esto al inicio de main.py (antes de la invocación):

from tools.web_search import configure_failures

configure_failures("academic", fail_rate=0.7)
configure_failures("news", fail_rate=0.3)

Ahora academic falla el 70% de las veces y news el 30%. El agente reintenta automáticamente.


Criterios de éxito

Tu proyecto está completo cuando cumples estos criterios:

  • Las búsquedas se ejecutan en paralelo — las 3 fuentes se lanzan simultáneamente (medible: el tiempo total es ~igual al de la fuente más lenta, no la suma)
  • Errores transitorios se recuperan — con configure_failures("academic", fail_rate=0.5), el agente reintenta y eventualmente obtiene resultados
  • Fallos permanentes no crashean — con configure_failures("academic", fail_rate=1.0), el agente continúa con web y news
  • El reporte incluye source availability — puedes ver cuántos intentos tuvo cada fuente y cuáles fallaron
  • Los resultados están deduplicados — no hay hallazgos repetidos de diferentes sub-queries que consultaron la misma fuente
  • La confianza refleja la availability — con todas las fuentes, confianza alta. Con fuentes faltantes, confianza ajustada a la baja
  • Los logs son estructurados — cada operación tiene un log JSON con request_id, nodo, duración, y status
  • El JSON del reporte v2 incluye version y source_availability — campos nuevos presentes y válidos

Escenarios de prueba

Test 1: Todas las fuentes OK (happy path)

# Sin configure_failures (default: 0% fail rate)
# python main.py "machine learning en medicina"
============================================================
  🔬 AI Research Assistant v2
  Tema: machine learning en medicina
  Request ID: a1b2c3d4
============================================================

📋 Paso 1: Descomponiendo tema en sub-queries...
   ✓ 4 sub-queries generadas

🔍 Paso 2: Buscando en 3 fuentes por sub-query (con retry)...
   Sub-query 1: 3 OK, 0 failed
   Sub-query 2: 3 OK, 0 failed
   Sub-query 3: 3 OK, 0 failed
   Sub-query 4: 3 OK, 0 failed
   Total: 12 exitosas, 0 fallidas de 12

🔀 Paso 3: Merge y deduplicación de resultados...
   ✓ 12 resultados raw → 12 únicos después de deduplicación

🧠 Paso 4: Sintetizando hallazgos...
   ✓ 4 hallazgos identificados

📝 Paso 5: Generando resumen ejecutivo...
   ✓ Resumen generado

📊 Paso 6: Calculando confianza del reporte...
   ✓ Confianza base: 82%, ajustada por availability: 82%

╔══════════════════════════════════════════════════════════╗
║       📄 REPORTE DE INVESTIGACIÓN (v2)                  ║
╚══════════════════════════════════════════════════════════╝

🔌 DISPONIBILIDAD DE FUENTES
────────────────────────────────────────────────────────────
  ✅ web: ok (305ms)
  ✅ academic: ok (420ms)
  ✅ news: ok (280ms)
  ...

Test 2: Una fuente falla intermitentemente (retry success)

from tools.web_search import configure_failures
configure_failures("academic", fail_rate=0.7)
🔍 Paso 2: Buscando en 3 fuentes por sub-query (con retry)...
   [Log: {"node": "search_academic", "event": "retry", "attempt": 1, "delay_s": 1.05, "error": "academic API: connection refused"}]
   [Log: {"node": "search_academic", "event": "retry", "attempt": 2, "delay_s": 2.13, "error": "academic API: 503 Service Unavailable"}]
   [Log: {"node": "search_academic", "event": "end", "duration_ms": 3850.2, "attempts": 3, "status": "success"}]
   Sub-query 1: 3 OK, 0 failed  ← academic se recuperó en el intento 3

🔌 DISPONIBILIDAD DE FUENTES
  ✅ web: ok (305ms)
  ✅ academic: ok (3 intentos, 3850ms)  ← tardó más pero funcionó
  ✅ news: ok (280ms)

Test 3: Una fuente falla permanentemente (graceful degradation)

configure_failures("academic", fail_rate=1.0)
🔍 Paso 2: Buscando en 3 fuentes por sub-query (con retry)...
   [Log: {"node": "search_academic", "event": "retry", "attempt": 1, ...}]
   [Log: {"node": "search_academic", "event": "retry", "attempt": 2, ...}]
   [Log: {"node": "search_academic", "event": "retry", "attempt": 3, ...}]
   [Log: {"node": "search_academic", "event": "error", "attempts": 4, "error": "academic API: 503"}]
   Sub-query 1: 2 OK, 1 failed  ← academic falló después de 4 intentos, web y news OK

📊 Paso 6: Calculando confianza del reporte...
   ✓ Confianza base: 78%, ajustada por availability: 72%  ← penalizada

🔌 DISPONIBILIDAD DE FUENTES
  ✅ web: ok (305ms)
  ❌ academic: failed (4 intentos, 7200ms) — academic API: 503 Service Unavailable
  ✅ news: ok (280ms)

Test 4: Todas las fuentes fallan

configure_failures("web", fail_rate=1.0)
configure_failures("academic", fail_rate=1.0)
configure_failures("news", fail_rate=1.0)
🔍 Paso 2: Buscando en 3 fuentes por sub-query (con retry)...
   Sub-query 1: 0 OK, 3 failed
   Total: 0 exitosas, 12 fallidas de 12

   ❌ Fuentes insuficientes (0 < 1). Abortando.

╔══════════════════════════════════════════════════════════╗
║           ❌ ERROR EN INVESTIGACIÓN                      ║
╚══════════════════════════════════════════════════════════╝

  Solo 0 fuentes respondieron. Mínimo requerido: 1
  Request ID: f4e5d6c7

  Disponibilidad de fuentes:
    ❌ web: failed (4 intentos)
    ❌ academic: failed (4 intentos)
    ❌ news: failed (4 intentos)

Test 5: Fallos intermitentes (stress test)

configure_failures("web", fail_rate=0.3)
configure_failures("academic", fail_rate=0.5)
configure_failures("news", fail_rate=0.2)

for topic in ["AI en educación", "quantum computing", "energías renovables"]:
    result = research_agent.invoke(topic, config={"configurable": {"thread_id": f"test-{topic}"}})
    ok = sum(1 for sa in result.get("source_availability", []) if sa["status"] == "ok")
    total = len(result.get("source_availability", []))
    print(f"  {topic}: {ok}/{total} fuentes OK, confianza {result.get('confidence', 'N/A')}")
# Output esperado (varía por random):
#   AI en educación: 10/12 fuentes OK, confianza 0.78
#   quantum computing: 11/12 fuentes OK, confianza 0.80
#   energías renovables: 9/12 fuentes OK, confianza 0.75

Errores comunes

1. ModuleNotFoundError: No module named 'utils'

Causa: El directorio utils/ no existe o no tiene __init__.py (aunque no siempre es necesario con sys.path.insert).

Solución: Ejecuta desde la raíz del proyecto y verifica la estructura:

cd research-assistant
ls utils/
# Debe mostrar: retry.py  logger.py
python main.py

2. El retry tarda demasiado en pruebas

Causa: Con base_delay=1.0 y 3 retries, cada búsqueda fallida tarda ~7 segundos (1s + 2s + 4s de backoff). Con 4 sub-queries y 3 fuentes, son potencialmente 12 × 7s = 84 segundos.

Solución: Para testing, reduce los delays:

from utils.retry import RetryConfig
test_config = RetryConfig(max_retries=2, base_delay=0.1, jitter=False)

O modifica RETRY_BASE_DELAY en config/settings.py para development.

3. La deduplicación no elimina hallazgos similares del LLM

Causa: merge_and_deduplicate compara por source_type + content[:100]. Si dos sub-queries distintas consultan la misma fuente, el content mock es diferente (porque incluye el query en el texto), así que no se detectan como duplicados.

Solución: La deduplicación actual es conservadora — solo elimina duplicados exactos. Para deduplicación semántica (hallazgos que dicen lo mismo con diferentes palabras), necesitarías embeddings o un LLM evaluador. Eso es materia del Módulo 11.

4. search_with_retry no usa la config de settings.py

Causa: La función search_with_retry crea su propio RetryConfig con valores hardcodeados.

Solución: Importa la configuración de settings:

from config.settings import RETRY_MAX_ATTEMPTS, RETRY_BASE_DELAY, RETRY_MAX_DELAY, RETRY_JITTER

retry_config = RetryConfig(
    max_retries=RETRY_MAX_ATTEMPTS,
    base_delay=RETRY_BASE_DELAY,
    max_delay=RETRY_MAX_DELAY,
    jitter=RETRY_JITTER,
)

5. Los logs de retry son excesivos en producción

Causa: Con muchas sub-queries y fuentes inestables, cada retry genera un log. 4 sub-queries × 3 fuentes × 3 retries = 36 logs de retry.

Solución: El ResearchLogger ya usa niveles — los retries son WARN. En producción, configura el nivel mínimo:

logging.basicConfig(level=logging.WARNING)  # Solo WARN y ERROR

6. SourceStatus no aparece en el reporte JSON

Causa: Estás usando el ResearchReport de v1 sin el campo source_availability.

Solución: Verifica que state/research_state.py tiene el modelo actualizado con SourceStatus y source_availability. El campo tiene default_factory=list, así que es compatible con v1 (aparece como lista vacía).

7. La confianza ajustada es siempre igual a la base

Causa: El availability_factor es total_ok / (total_ok + total_failed). Si no hay fallos, availability_factor = 1.0 y la fórmula base * (0.7 + 0.3 * 1.0) = base * 1.0.

Solución: Esto es correcto — la confianza solo se penaliza cuando hay fuentes faltantes. Con todas las fuentes disponibles, la confianza ajustada es igual a la base.

8. TypeError al pasar logger como argumento de @task

Causa: El logger no es JSON-serializable, y @task intenta serializarlo para el checkpoint.

Solución: Usa un logger global en vez de pasarlo como argumento, o exclúyelo de la serialización. En el código de este proyecto, search_with_retry recibe el logger como argumento — asegúrate de que tu checkpointer puede manejar esto, o usa el pattern de logger global mostrado en la cápsula 07.


Lo que viene: Módulo 8 — Memoria y Persistencia

Tu Research Assistant v2 es robusto: reintenta errores, busca en paralelo, degrada gracefully, y logea todo. Pero tiene una limitación fundamental: cada investigación empieza de cero.

Si investigaste "inteligencia artificial en medicina" ayer, y hoy investigas "IA aplicada a diagnóstico", el agente no sabe que ya tiene información relevante. Si el proceso se interrumpe a mitad de una investigación larga, pierdes todo el progreso.

El Módulo 8 agrega dos capacidades:

  • Checkpointing persistente con PostgresSaver — el agente guarda su estado en cada paso. Si se cae en el paso 4, resume desde el paso 3 sin repetir trabajo
  • Long-term memory — el agente recuerda investigaciones anteriores, preferencias del usuario, y puede reutilizar hallazgos previos relevantes

La v3 no es solo más robusta — es más inteligente. Un agente que recuerda es fundamentalmente diferente de uno que no.


Recursos del proyecto

  1. LangGraph Functional API — Referencia de @entrypoint y @task, incluyendo Futures
  2. Exponential Backoff and Jitter — AWS blog sobre por qué jitter es esencial en retry logic
  3. Python logging Cookbook — Logging structured con JSON formatters
  4. Pydantic v2 — Field Defaults — Cómo hacer modelos backwards-compatible con default_factory
  5. LangGraph Persistence — Checkpointing y MemorySaver
  6. Graceful Degradation in Distributed Systems — Capítulo 8 de "Designing Data-Intensive Applications" (Kleppmann)

Módulo 7 — LangChain & LangGraph: From Chains to Agents