Módulo 7: Casos de Uso

8. Proyecto: Use Case Selector

Descripción

En este proyecto construyes un Use Case Selector: un router inteligente que recibe un input (PDF, imagen, audio, video, texto), detecta automáticamente su tipo, y aplica el pipeline correcto. Es el "cerebro" que conecta todos los patrones de diseño que aprendiste en este módulo en un solo sistema.

Por qué importa: En un sistema multimodal real, el usuario no dice "aplica el pipeline de Document Q&A". Sube un archivo y espera que el sistema haga lo correcto. El Use Case Selector es esa capa de inteligencia: recibe cualquier input, decide qué hacer con él, y ejecuta el pipeline apropiado. Es el componente más crítico del Document Analyzer que construirás en el Módulo 8.

Lo que vas a construir:

┌──────────────┐     ┌──────────────┐     ┌──────────────┐     ┌──────────────┐
│  Input       │────▶│  Detector    │────▶│  Pipeline    │────▶│  Router      │
│  (cualquier) │     │  de tipo     │     │  Registry    │     │  (ejecutar)  │
└──────────────┘     └──────────────┘     └──────────────┘     └──────┬───────┘
                                                                       │
                                                                       ▼
                                                               ┌──────────────┐
                                                               │  Resultado   │
                                                               │  + tracking  │
                                                               └──────────────┘

Especificaciones

Input

  • Ruta a archivo — PDF, imagen, audio, video, texto
  • Pregunta opcional — Si el usuario quiere Q&A sobre el archivo
  • Configuración — Modelo preferido, presupuesto máximo, calidad deseada

Output

  • Tipo detectado — Qué tipo de input es
  • Pipeline seleccionado — Qué pipeline se aplicó
  • Resultado — Output del pipeline ejecutado
  • Metadatos — Costo, duración, modelo usado

Componentes

#ComponenteResponsabilidad
1ConfigConfiguración global y por pipeline
2Type DetectorDetectar tipo de input
3Pipeline RegistryRegistrar y gestionar pipelines disponibles
4RouterSeleccionar y ejecutar pipeline
5Cost TrackerRastrear costos por operación
6MainOrquestar todo

Paso 1: Configuración

from dataclasses import dataclass, field

@dataclass
class PipelineConfig:
    model: str = "gpt-4o-mini"
    max_tokens: int = 500
    temperature: float = 0.2
    max_cost_per_request: float = 0.50
    enable_cache: bool = True
    enable_logging: bool = True

@dataclass
class UseCaseSelectorConfig:
    default_pipeline_config: PipelineConfig = field(default_factory=PipelineConfig)
    supported_types: set = field(default_factory=lambda: {
        "document", "image", "audio", "video", "text"
    })
    max_file_size_mb: float = 100.0
    fallback_enabled: bool = True

    type_extensions: dict = field(default_factory=lambda: {
        ".pdf": "document",
        ".doc": "document",
        ".docx": "document",
        ".png": "image",
        ".jpg": "image",
        ".jpeg": "image",
        ".gif": "image",
        ".webp": "image",
        ".bmp": "image",
        ".mp3": "audio",
        ".wav": "audio",
        ".m4a": "audio",
        ".ogg": "audio",
        ".flac": "audio",
        ".mp4": "video",
        ".avi": "video",
        ".mov": "video",
        ".mkv": "video",
        ".webm": "video",
        ".txt": "text",
        ".md": "text",
        ".csv": "text",
        ".json": "text",
    })

Paso 2: Detector de Tipo

from pathlib import Path
import mimetypes

class TypeDetector:
    def __init__(self, config: UseCaseSelectorConfig):
        self.config = config

    def detect(self, path: str) -> dict:
        p = Path(path)

        if not p.exists():
            return {
                "type": "unknown",
                "error": "file_not_found",
                "path": path
            }

        ext = p.suffix.lower()
        size_mb = p.stat().st_size / (1024 * 1024)

        if size_mb > self.config.max_file_size_mb:
            return {
                "type": "unknown",
                "error": "file_too_large",
                "size_mb": round(size_mb, 2),
                "max_mb": self.config.max_file_size_mb,
                "path": path
            }

        detected_type = self.config.type_extensions.get(ext)

        if not detected_type:
            mime_type, _ = mimetypes.guess_type(path)
            if mime_type:
                detected_type = self._mime_to_type(mime_type)

        if not detected_type:
            detected_type = "unknown"

        return {
            "type": detected_type,
            "extension": ext,
            "size_mb": round(size_mb, 2),
            "path": path,
            "mime_type": mimetypes.guess_type(path)[0]
        }

    def _mime_to_type(self, mime: str) -> str | None:
        if mime.startswith("image/"):
            return "image"
        if mime.startswith("audio/"):
            return "audio"
        if mime.startswith("video/"):
            return "video"
        if mime.startswith("text/"):
            return "text"
        if mime == "application/pdf":
            return "document"
        return None

    def detect_multiple(self, paths: list[str]) -> list[dict]:
        return [self.detect(path) for path in paths]

    def validate(self, path: str) -> dict:
        detection = self.detect(path)

        if detection["type"] == "unknown":
            return {
                "valid": False,
                "detection": detection,
                "message": detection.get("error", "Tipo no soportado")
            }

        if detection["type"] not in self.config.supported_types:
            return {
                "valid": False,
                "detection": detection,
                "message": f"Tipo '{detection['type']}' no está habilitado"
            }

        return {
            "valid": True,
            "detection": detection
        }

Paso 3: Pipeline Registry

from typing import Callable, Any

class Pipeline:
    def __init__(
        self,
        name: str,
        input_type: str,
        handler: Callable,
        description: str = "",
        supports_question: bool = False,
        estimated_cost_per_call: float = 0.01
    ):
        self.name = name
        self.input_type = input_type
        self.handler = handler
        self.description = description
        self.supports_question = supports_question
        self.estimated_cost_per_call = estimated_cost_per_call


class PipelineRegistry:
    def __init__(self):
        self.pipelines: dict[str, list[Pipeline]] = {}

    def register(self, pipeline: Pipeline) -> None:
        if pipeline.input_type not in self.pipelines:
            self.pipelines[pipeline.input_type] = []
        self.pipelines[pipeline.input_type].append(pipeline)

    def get_pipeline(self, input_type: str, has_question: bool = False) -> Pipeline | None:
        candidates = self.pipelines.get(input_type, [])

        if not candidates:
            return None

        if has_question:
            qa_pipelines = [p for p in candidates if p.supports_question]
            if qa_pipelines:
                return qa_pipelines[0]

        return candidates[0]

    def list_pipelines(self) -> dict:
        result = {}
        for input_type, pipelines in self.pipelines.items():
            result[input_type] = [
                {
                    "name": p.name,
                    "description": p.description,
                    "supports_question": p.supports_question,
                    "estimated_cost": p.estimated_cost_per_call
                }
                for p in pipelines
            ]
        return result

    def get_all_supported_types(self) -> set[str]:
        return set(self.pipelines.keys())

Registrar pipelines

from openai import OpenAI
import base64
import fitz

client = OpenAI()

def document_extract_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
    config = config or PipelineConfig()
    doc = fitz.open(path)
    pages = []
    for i in range(len(doc)):
        text = doc[i].get_text().strip()
        if text:
            pages.append({"page": i + 1, "text": text})
    doc.close()

    full_text = "\n\n".join(p["text"] for p in pages)

    if question:
        response = client.chat.completions.create(
            model=config.model,
            messages=[
                {
                    "role": "system",
                    "content": "Responde SOLO con base en el contexto. Si no encuentras la respuesta, dilo. Cita la página."
                },
                {
                    "role": "user",
                    "content": f"Documento:\n{full_text[:8000]}\n\nPregunta: {question}"
                }
            ],
            max_tokens=config.max_tokens,
            temperature=config.temperature
        )
        return {
            "type": "document_qa",
            "answer": response.choices[0].message.content,
            "pages_processed": len(pages),
            "model": config.model
        }

    return {
        "type": "document_extraction",
        "pages": pages,
        "total_pages": len(pages),
        "total_characters": len(full_text)
    }


def image_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
    config = config or PipelineConfig()
    with open(path, "rb") as f:
        b64 = base64.b64encode(f.read()).decode()

    prompt = question or "Analiza esta imagen. Describe qué contiene, extrae cualquier texto visible, y clasifica el tipo de contenido."

    response = client.chat.completions.create(
        model=config.model,
        messages=[{
            "role": "user",
            "content": [
                {"type": "text", "text": prompt},
                {"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{b64}"}}
            ]
        }],
        max_tokens=config.max_tokens,
        temperature=config.temperature
    )

    return {
        "type": "image_analysis",
        "analysis": response.choices[0].message.content,
        "model": config.model
    }


def audio_transcribe_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
    config = config or PipelineConfig()
    with open(path, "rb") as f:
        transcript = client.audio.transcriptions.create(
            model="whisper-1",
            file=f,
            response_format="text"
        )

    if question:
        response = client.chat.completions.create(
            model=config.model,
            messages=[
                {
                    "role": "system",
                    "content": "Responde basándote en la transcripción del audio."
                },
                {
                    "role": "user",
                    "content": f"Transcripción:\n{transcript[:6000]}\n\nPregunta: {question}"
                }
            ],
            max_tokens=config.max_tokens,
            temperature=config.temperature
        )
        return {
            "type": "audio_qa",
            "answer": response.choices[0].message.content,
            "transcript": transcript,
            "model": config.model
        }

    return {
        "type": "audio_transcription",
        "transcript": transcript,
        "transcript_length": len(transcript)
    }


def video_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
    config = config or PipelineConfig()
    import cv2
    import os

    cap = cv2.VideoCapture(path)
    fps = cap.get(cv2.CAP_PROP_FPS)
    total_frames = int(cap.get(cv2.CAP_PROP_FRAME_COUNT))
    duration = total_frames / fps if fps > 0 else 0

    target_frames = min(15, max(5, int(duration / 10)))
    interval = max(1, int(total_frames / target_frames))

    os.makedirs("/tmp/uc_selector_frames", exist_ok=True)
    descriptions = []
    frame_id = 0
    captured = 0

    while True:
        ret, frame = cap.read()
        if not ret:
            break
        if frame_id % interval == 0 and captured < target_frames:
            frame_path = f"/tmp/uc_selector_frames/frame_{frame_id}.jpg"
            cv2.imwrite(frame_path, frame)

            with open(frame_path, "rb") as f:
                b64 = base64.b64encode(f.read()).decode()

            ts = frame_id / fps
            minutes = int(ts // 60)
            seconds = int(ts % 60)

            response = client.chat.completions.create(
                model="gpt-4o-mini",
                messages=[{
                    "role": "user",
                    "content": [
                        {"type": "text", "text": "Describe esta escena en 1-2 oraciones."},
                        {"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{b64}"}}
                    ]
                }],
                max_tokens=100
            )
            descriptions.append(f"{minutes}:{seconds:02d} - {response.choices[0].message.content}")
            captured += 1

        frame_id += 1

    cap.release()

    timeline = "\n".join(descriptions)

    summary_prompt = f"Análisis de frames del video:\n\n{timeline}\n\n"
    if question:
        summary_prompt += f"Pregunta: {question}"
    else:
        summary_prompt += "Genera un resumen del video con los puntos clave."

    response = client.chat.completions.create(
        model=config.model,
        messages=[{"role": "user", "content": summary_prompt}],
        max_tokens=config.max_tokens
    )

    return {
        "type": "video_analysis",
        "summary": response.choices[0].message.content,
        "frames_analyzed": captured,
        "duration_seconds": round(duration, 2),
        "model": config.model
    }


def text_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
    config = config or PipelineConfig()
    with open(path, "r", encoding="utf-8") as f:
        text = f.read()

    prompt = question or "Analiza este texto. Resume los puntos clave y extrae información relevante."

    response = client.chat.completions.create(
        model=config.model,
        messages=[
            {
                "role": "user",
                "content": f"Texto:\n{text[:8000]}\n\n{prompt}"
            }
        ],
        max_tokens=config.max_tokens,
        temperature=config.temperature
    )

    return {
        "type": "text_analysis",
        "result": response.choices[0].message.content,
        "text_length": len(text),
        "model": config.model
    }


def create_default_registry() -> PipelineRegistry:
    registry = PipelineRegistry()

    registry.register(Pipeline(
        name="document_extractor",
        input_type="document",
        handler=document_extract_handler,
        description="Extrae texto de PDFs y responde preguntas sobre documentos",
        supports_question=True,
        estimated_cost_per_call=0.01
    ))

    registry.register(Pipeline(
        name="image_analyzer",
        input_type="image",
        handler=image_analyze_handler,
        description="Analiza imágenes: describe contenido, extrae texto, clasifica",
        supports_question=True,
        estimated_cost_per_call=0.005
    ))

    registry.register(Pipeline(
        name="audio_transcriber",
        input_type="audio",
        handler=audio_transcribe_handler,
        description="Transcribe audio y responde preguntas sobre el contenido",
        supports_question=True,
        estimated_cost_per_call=0.02
    ))

    registry.register(Pipeline(
        name="video_analyzer",
        input_type="video",
        handler=video_analyze_handler,
        description="Analiza video extrayendo frames y generando resumen",
        supports_question=True,
        estimated_cost_per_call=0.10
    ))

    registry.register(Pipeline(
        name="text_analyzer",
        input_type="text",
        handler=text_analyze_handler,
        description="Analiza archivos de texto y responde preguntas",
        supports_question=True,
        estimated_cost_per_call=0.005
    ))

    return registry

Paso 4: Router

import time
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("use_case_selector")

class UseCaseRouter:
    def __init__(
        self,
        config: UseCaseSelectorConfig = None,
        registry: PipelineRegistry = None,
        cost_tracker = None
    ):
        self.config = config or UseCaseSelectorConfig()
        self.detector = TypeDetector(self.config)
        self.registry = registry or create_default_registry()
        self.cost_tracker = cost_tracker
        self.history: list[dict] = []

    def route(self, path: str, question: str = None) -> dict:
        start_time = time.time()

        validation = self.detector.validate(path)
        if not validation["valid"]:
            return {
                "status": "error",
                "error": validation["message"],
                "detection": validation["detection"]
            }

        detection = validation["detection"]
        input_type = detection["type"]
        has_question = question is not None and question.strip() != ""

        pipeline = self.registry.get_pipeline(input_type, has_question)
        if not pipeline:
            return {
                "status": "error",
                "error": f"No hay pipeline registrado para tipo '{input_type}'",
                "detection": detection
            }

        logger.info(f"Routing: {path}{pipeline.name} (type={input_type}, question={has_question})")

        try:
            result = pipeline.handler(
                path,
                question=question,
                config=self.config.default_pipeline_config
            )

            duration_ms = round((time.time() - start_time) * 1000, 2)

            entry = {
                "status": "success",
                "detection": detection,
                "pipeline": pipeline.name,
                "result": result,
                "duration_ms": duration_ms,
                "estimated_cost": pipeline.estimated_cost_per_call,
                "question": question
            }

            if self.cost_tracker:
                self.cost_tracker.track_chat(
                    self.config.default_pipeline_config.model,
                    input_tokens=1000,
                    output_tokens=500
                )

            self.history.append(entry)
            logger.info(f"Success: {pipeline.name} in {duration_ms}ms")

            return entry

        except Exception as e:
            duration_ms = round((time.time() - start_time) * 1000, 2)

            error_entry = {
                "status": "error",
                "detection": detection,
                "pipeline": pipeline.name,
                "error": str(e),
                "error_type": type(e).__name__,
                "duration_ms": duration_ms
            }

            self.history.append(error_entry)
            logger.error(f"Error in {pipeline.name}: {e}")

            if self.config.fallback_enabled:
                return self._try_fallback(path, question, error_entry)

            return error_entry

    def _try_fallback(self, path: str, question: str, original_error: dict) -> dict:
        logger.info("Attempting fallback...")

        try:
            with open(path, "rb") as f:
                content = f.read()

            if len(content) < 10000:
                text_content = content.decode("utf-8", errors="replace")
                prompt = f"Analiza este contenido:\n{text_content[:6000]}"
                if question:
                    prompt += f"\n\nPregunta: {question}"

                response = client.chat.completions.create(
                    model="gpt-4o-mini",
                    messages=[{"role": "user", "content": prompt}],
                    max_tokens=500
                )

                return {
                    "status": "fallback",
                    "original_error": original_error,
                    "result": {
                        "type": "fallback_analysis",
                        "analysis": response.choices[0].message.content
                    }
                }
        except Exception:
            pass

        return original_error

    def route_multiple(self, paths: list[str], question: str = None) -> list[dict]:
        return [self.route(path, question) for path in paths]

    def get_stats(self) -> dict:
        total = len(self.history)
        successes = sum(1 for h in self.history if h["status"] == "success")
        errors = sum(1 for h in self.history if h["status"] == "error")
        fallbacks = sum(1 for h in self.history if h["status"] == "fallback")

        durations = [h["duration_ms"] for h in self.history if "duration_ms" in h]
        total_cost = sum(h.get("estimated_cost", 0) for h in self.history)

        pipeline_counts = {}
        for h in self.history:
            p = h.get("pipeline", "unknown")
            pipeline_counts[p] = pipeline_counts.get(p, 0) + 1

        return {
            "total_requests": total,
            "successes": successes,
            "errors": errors,
            "fallbacks": fallbacks,
            "success_rate": round(successes / total * 100, 1) if total > 0 else 0,
            "avg_duration_ms": round(sum(durations) / len(durations), 2) if durations else 0,
            "total_estimated_cost": round(total_cost, 4),
            "by_pipeline": pipeline_counts
        }

Paso 5: Cost Tracker Integrado

class UseCaseCostTracker:
    def __init__(self):
        self.records: list[dict] = []

    def track(self, pipeline_name: str, input_type: str, cost: float, duration_ms: float):
        self.records.append({
            "pipeline": pipeline_name,
            "input_type": input_type,
            "cost": cost,
            "duration_ms": duration_ms,
            "timestamp": time.time()
        })

    def summary(self) -> dict:
        total_cost = sum(r["cost"] for r in self.records)
        by_pipeline = {}
        for r in self.records:
            p = r["pipeline"]
            if p not in by_pipeline:
                by_pipeline[p] = {"cost": 0, "count": 0}
            by_pipeline[p]["cost"] += r["cost"]
            by_pipeline[p]["count"] += 1

        return {
            "total_cost": round(total_cost, 4),
            "total_requests": len(self.records),
            "by_pipeline": {
                k: {"cost": round(v["cost"], 4), "count": v["count"]}
                for k, v in by_pipeline.items()
            }
        }

    def check_budget(self, max_daily_cost: float) -> dict:
        today_records = [
            r for r in self.records
            if r["timestamp"] > time.time() - 86400
        ]
        today_cost = sum(r["cost"] for r in today_records)

        return {
            "today_cost": round(today_cost, 4),
            "daily_budget": max_daily_cost,
            "remaining": round(max_daily_cost - today_cost, 4),
            "within_budget": today_cost < max_daily_cost,
            "usage_percent": round(today_cost / max_daily_cost * 100, 1) if max_daily_cost > 0 else 0
        }

Paso 6: Main — Orquestación Completa

class UseCaseSelector:
    def __init__(self, config: UseCaseSelectorConfig = None):
        self.config = config or UseCaseSelectorConfig()
        self.registry = create_default_registry()
        self.cost_tracker = UseCaseCostTracker()
        self.router = UseCaseRouter(
            config=self.config,
            registry=self.registry,
            cost_tracker=self.cost_tracker
        )

    def process(self, path: str, question: str = None) -> dict:
        return self.router.route(path, question)

    def process_batch(self, paths: list[str], question: str = None) -> list[dict]:
        return self.router.route_multiple(paths, question)

    def info(self) -> dict:
        return {
            "supported_types": list(self.config.supported_types),
            "available_pipelines": self.registry.list_pipelines(),
            "config": {
                "model": self.config.default_pipeline_config.model,
                "max_file_size_mb": self.config.max_file_size_mb,
                "fallback_enabled": self.config.fallback_enabled,
                "cache_enabled": self.config.default_pipeline_config.enable_cache
            }
        }

    def stats(self) -> dict:
        return {
            "routing": self.router.get_stats(),
            "costs": self.cost_tracker.summary()
        }

    def add_pipeline(self, pipeline: Pipeline) -> None:
        self.registry.register(pipeline)

Uso completo

selector = UseCaseSelector()

print(selector.info())

result = selector.process("contrato.pdf")
print(f"Tipo: {result['detection']['type']}")
print(f"Pipeline: {result['pipeline']}")
print(f"Resultado: {result['result']}")

result = selector.process(
    "contrato.pdf",
    question="¿Cuál es la cláusula de penalización?"
)
print(f"Respuesta: {result['result']['answer']}")

result = selector.process("producto.jpg")
print(f"Análisis: {result['result']['analysis']}")

result = selector.process("reunion.mp3")
print(f"Transcripción: {result['result']['transcript'][:200]}")

result = selector.process("clase.mp4")
print(f"Resumen: {result['result']['summary']}")

print(selector.stats())

Extensión 1: Auto-Pipeline Selection con LLM

En lugar de solo usar la extensión del archivo, usar un LLM para decidir qué pipeline aplicar basándose en el contenido:

def smart_pipeline_selection(
    path: str,
    question: str,
    detection: dict,
    available_pipelines: dict
) -> str:
    pipelines_desc = "\n".join(
        f"- {p['name']}: {p['description']}"
        for type_pipelines in available_pipelines.values()
        for p in type_pipelines
    )

    response = client.chat.completions.create(
        model="gpt-4o-mini",
        messages=[{
            "role": "user",
            "content": (
                f"Archivo: {path}\n"
                f"Tipo detectado: {detection['type']}\n"
                f"Extensión: {detection['extension']}\n"
                f"Tamaño: {detection['size_mb']}MB\n"
                f"Pregunta del usuario: {question or 'Ninguna'}\n\n"
                f"Pipelines disponibles:\n{pipelines_desc}\n\n"
                "¿Qué pipeline es el más apropiado? Responde SOLO con el nombre del pipeline."
            )
        }],
        max_tokens=30,
        temperature=0
    )
    return response.choices[0].message.content.strip()

Extensión 2: Batch Routing con Prioridades

Para procesar múltiples archivos con diferentes prioridades:

from dataclasses import dataclass
from typing import Optional

@dataclass
class RoutingRequest:
    path: str
    question: Optional[str] = None
    priority: int = 5

class BatchRouter:
    def __init__(self, selector: UseCaseSelector):
        self.selector = selector
        self.queue: list[RoutingRequest] = []

    def add(self, path: str, question: str = None, priority: int = 5):
        self.queue.append(RoutingRequest(path=path, question=question, priority=priority))

    def process_all(self) -> list[dict]:
        sorted_queue = sorted(self.queue, key=lambda r: r.priority)

        results = []
        for request in sorted_queue:
            result = self.selector.process(request.path, request.question)
            result["priority"] = request.priority
            results.append(result)

        self.queue.clear()
        return results

    def process_by_type(self) -> dict[str, list[dict]]:
        by_type: dict[str, list[RoutingRequest]] = {}
        for req in self.queue:
            detection = self.selector.router.detector.detect(req.path)
            input_type = detection["type"]
            if input_type not in by_type:
                by_type[input_type] = []
            by_type[input_type].append(req)

        results = {}
        for input_type, requests in by_type.items():
            results[input_type] = []
            for req in sorted(requests, key=lambda r: r.priority):
                result = self.selector.process(req.path, req.question)
                results[input_type].append(result)

        self.queue.clear()
        return results

Uso:

batch = BatchRouter(selector)

batch.add("urgente.pdf", priority=1)
batch.add("producto1.jpg", priority=5)
batch.add("producto2.jpg", priority=5)
batch.add("reunion.mp3", question="¿Qué decisiones se tomaron?", priority=3)

results = batch.process_all()
for r in results:
    print(f"[P{r['priority']}] {r['pipeline']}: {r['status']}")

Troubleshooting

Problema 1: Tipo de archivo no detectado

Síntoma: TypeDetector.detect retorna "unknown" para un archivo válido.

Solución: Agregar la extensión al config:

config = UseCaseSelectorConfig()
config.type_extensions[".xlsx"] = "document"
config.type_extensions[".pptx"] = "document"

Problema 2: Pipeline incorrecto seleccionado

Síntoma: Un PDF escaneado se enruta a document_extractor pero no hay texto que extraer.

Solución: Agregar detección de contenido después del tipo:

def detect_pdf_subtype(path: str) -> str:
    import fitz
    doc = fitz.open(path)
    total_text = sum(len(doc[i].get_text().strip()) for i in range(len(doc)))
    total_images = sum(len(doc[i].get_images()) for i in range(len(doc)))
    doc.close()

    if total_text < 100 and total_images > 0:
        return "scanned_document"
    return "text_document"

Problema 3: Videos largos consumen demasiado

Síntoma: Un video de 2 horas genera 100+ llamadas a Vision API.

Solución: El video_analyze_handler ya limita frames con target_frames = min(15, ...). Para más control:

config = PipelineConfig()
config.max_tokens = 300

selector = UseCaseSelector(UseCaseSelectorConfig(
    default_pipeline_config=config,
    max_file_size_mb=50.0
))

Problema 4: Costo acumulado sin control

Síntoma: El cost tracker muestra gasto alto pero no hay alertas.

Solución:

budget = selector.cost_tracker.check_budget(max_daily_cost=5.0)
if not budget["within_budget"]:
    print(f"ALERTA: Presupuesto excedido. Gastado: ${budget['today_cost']}")

Checklist de Completitud

Funcionalidad core

  • Detecta PDF y enruta a document pipeline
  • Detecta imagen y enruta a image pipeline
  • Detecta audio y enruta a audio pipeline
  • Detecta video y enruta a video pipeline
  • Detecta texto y enruta a text pipeline
  • Maneja archivos no soportados con error claro
  • Soporta pregunta opcional para Q&A

Robustez

  • Valida que el archivo existe
  • Verifica tamaño máximo
  • Maneja errores de API con fallback
  • Registra historial de operaciones
  • Trackea costos estimados

Extensibilidad

  • Se pueden registrar pipelines nuevos
  • La configuración es flexible
  • Soporta batch processing

Ejercicios

Ejercicio 1: Pipeline personalizado

Agrega un pipeline para archivos CSV que lea el archivo, analice las columnas con un LLM, y genere un resumen estadístico.

Ver solución
import csv

def csv_analyze_handler(path: str, question: str = None, config: PipelineConfig = None) -> dict:
    config = config or PipelineConfig()

    with open(path, "r", encoding="utf-8") as f:
        reader = csv.reader(f)
        headers = next(reader)
        rows = list(reader)

    sample = rows[:20]
    sample_text = "\n".join([",".join(headers)] + [",".join(row) for row in sample])

    prompt = (
        f"Archivo CSV con {len(rows)} filas y {len(headers)} columnas.\n"
        f"Columnas: {', '.join(headers)}\n\n"
        f"Muestra de datos:\n{sample_text}\n\n"
    )

    if question:
        prompt += f"Pregunta: {question}"
    else:
        prompt += "Analiza estos datos: describe las columnas, identifica patrones, y genera un resumen."

    response = client.chat.completions.create(
        model=config.model,
        messages=[{"role": "user", "content": prompt}],
        max_tokens=config.max_tokens
    )

    return {
        "type": "csv_analysis",
        "result": response.choices[0].message.content,
        "rows": len(rows),
        "columns": headers
    }

selector.add_pipeline(Pipeline(
    name="csv_analyzer",
    input_type="text",
    handler=csv_analyze_handler,
    description="Analiza archivos CSV",
    supports_question=True,
    estimated_cost_per_call=0.005
))

Ejercicio 2: Multi-file routing

Extiende el selector para que reciba múltiples archivos y los combine en un solo análisis cuando sean del mismo tema.

Ver solución
def process_related_files(
    selector: UseCaseSelector,
    paths: list[str],
    question: str = None
) -> dict:
    individual_results = []
    for path in paths:
        result = selector.process(path)
        if result["status"] == "success":
            individual_results.append({
                "path": path,
                "type": result["detection"]["type"],
                "summary": str(result["result"])[:500]
            })

    if len(individual_results) > 1:
        summaries = "\n\n".join(
            f"[{r['type']}] {r['path']}:\n{r['summary']}"
            for r in individual_results
        )

        prompt = f"Análisis de {len(individual_results)} archivos:\n\n{summaries}\n\n"
        if question:
            prompt += f"Pregunta: {question}"
        else:
            prompt += "Genera un análisis integrado que conecte la información de todos los archivos."

        response = client.chat.completions.create(
            model="gpt-4o-mini",
            messages=[{"role": "user", "content": prompt}],
            max_tokens=600
        )

        return {
            "type": "multi_file_analysis",
            "integrated_result": response.choices[0].message.content,
            "individual_results": individual_results,
            "files_processed": len(individual_results)
        }

    return {
        "type": "single_file",
        "results": individual_results
    }

Resumen

El Use Case Selector es:

  • Un router que detecta tipo de input y selecciona el pipeline correcto
  • Un registry que gestiona pipelines disponibles y permite agregar nuevos
  • Un tracker que monitorea costos y rendimiento
  • La base del Document Analyzer del Módulo 8

Componentes implementados:

ComponenteClaseResponsabilidad
ConfigUseCaseSelectorConfig, PipelineConfigConfiguración flexible
DetectorTypeDetectorDetectar tipo de archivo
RegistryPipelineRegistryGestionar pipelines
RouterUseCaseRouterSeleccionar y ejecutar
CostUseCaseCostTrackerRastrear costos
MainUseCaseSelectorOrquestar todo

Próximo módulo: Módulo 8 — Proyecto Final: Document Analyzer Multimodal. El Use Case Selector que construiste aquí se integra como componente de routing en un sistema completo que procesa documentos, extrae datos con vision, indexa para Q&A con RAG, y genera audio con TTS.