> For the complete documentation index, see [llms.txt](https://public-intelligence.gitbook.io/taina-agente-ia-ogtic/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://public-intelligence.gitbook.io/taina-agente-ia-ogtic/arquitectura-y-conceptos/estructura-proyecto/sistema-ingesta.md).

# Sistema de ingesta

El sistema de ingesta de Taina está diseñado para procesar y almacenar documentos de servicios gubernamentales en una base de datos vectorial (ChromaDB) para habilitar búsquedas semánticas y RAG (Retrieval Augmented Generation).

## Arquitectura del sistema de ingesta

### Componentes principales

```
src/
├── ingest.py                    # Punto de entrada principal
├── embeddings/
│   └── gemini_embedding.py     # Embeddings de Google Gemini
├── vectors/
│   └── chroma_vector.py        # Integración con ChromaDB
└── utils/
    ├── doc_loader.py           # Carga de documentos
    ├── doc_split.py            # División de documentos
    ├── logger.py               # Sistema de logging
    └── rate_limiter.py         # Control de rate limiting
```

## Flujo de procesamiento

### 1. Punto de entrada: `ingest.py`

**Función principal**: `ingest()`

* **Parámetros**:
  * `json_dir`: Directorio con JSONs
  * `persist_path`: Ruta de ChromaDB
  * `collection_name`: Nombre de colección
  * `embeddings_model`: Modelo de embeddings

**Flujo**:

```python
def ingest(page_urls, docs_dir, json_dir, metadata_keys, persist_path, collection_name, embeddings_model):
    # 1. Obtener embeddings
    embeddings = embeddings_model or get_google_embeddings()
    
    # 2. Crear directorio de persistencia
    persist_path.mkdir(parents=True, exist_ok=True)
    
    # 3. Cargar documentos desde fuente json
    documents = []
    
    # ... (otras líneas de código)
    
    if json_dir:
        json_docs = load_json_from_directory(docs_dir=json_dir, metadata_keys=metadata_keys)
        documents.extend(json_docs)
    
    # 4. Ingestar en ChromaDB
    get_chroma_ingest(docs=documents, directory=persist_path, collection_name=collection_name, embeddings=embeddings)
```

### 2. Carga de documentos: `doc_loader.py`

#### `load_json_from_directory(docs_dir: Path, metadata_keys: List[str]) -> List[Document]`

* **Propósito**: Cargar y procesar JSONs de servicios gubernamentales
* **Funcionalidad avanzada**: Creación de chunks granulares para RAG
* **Tipos de documentos generados**:
  1. **Servicio Principal**: Información general
  2. **Variaciones**: Procedimientos específicos
  3. **Ubicaciones**: Oficinas y horarios
  4. **Canales Digitales**: Enlaces online
  5. **Contacto**: Información de contacto

**Estructura de chunks JSON**:

```python
# Chunk 1: Servicio Principal
main_content = f"""Servicio: {datos.get('nombre', 'N/A')}
Institución: {datos.get('institucion_responsable', 'N/A')} ({datos.get('sigla_de_institucion', 'N/A')})

Descripción:
{datos.get('descripcion', 'N/A')}

Objetivo:
{datos.get('objetivo', 'N/A')}

Canales disponibles: {', '.join(datos.get('canales_disponibles', []))}"""

# Chunk 2-N: Variaciones del Servicio
for i, variacion in enumerate(variaciones):
    var_content = f"""Servicio: {datos.get('nombre', 'N/A')}
Variación: {variacion.get('nombre', 'N/A')}
Precio: {variacion.get('moneda', 'RD$')} {variacion.get('precio', 'Gratuito')}
Procedimiento: {variacion.get('procedimiento_variacion_servicio', 'N/A')}
Requisitos: {chr(10).join(f"• {req}" for req in requisitos_list)}
Documentos necesarios: {chr(10).join(f"• {doc}" for doc in documentos_list)}"""

# Chunk N+1: Ubicaciones
ubicaciones_content = "Oficinas y ubicaciones disponibles:\n"
for oficina in oficinas:
    ubicaciones_content += f"\n• {oficina['nombre']}\n"
    ubicaciones_content += f"  Dirección: {oficina['calle']} #{oficina['numero']}\n"
    ubicaciones_content += f"  {oficina.get('municipio', '')}, {oficina.get('provincia', '')}\n"

# Chunk N+2: Canal Digital
digital_content = f"""Servicio: {datos.get('nombre', 'N/A')}
Institución: {datos.get('sigla_de_institucion', 'N/A')}
Canal Digital: {datos.get('canal_digital', 'Portal Web')}
Enlace: {datos.get('enlace_canal_digital', 'N/A')}"""

# Chunk N+3: Contacto
contacto_content = f"""Información de contacto - {datos.get('nombre', 'Servicio')}
Institución: {oficina.get('nombre', 'N/A')}
Portal: {oficina.get('portal_institucional', 'N/A')}
Dirección: {oficina.get('direccion_oficina_principal', 'N/A')}
Teléfonos: {', '.join(telefonos)}
Correos: {', '.join(emails)}"""
```

### 3. División de documentos: `doc_split.py`

**Función**: `split_documents(documents)`

* **Tecnología**: `RecursiveCharacterTextSplitter`
* **Configuración**:
  * `chunk_size=1000`: Tamaño de cada chunk
  * `chunk_overlap=200`: Solapamiento entre chunks
  * `add_start_index=True`: Índices de inicio

### 4. Embeddings: `gemini_embedding.py`

#### `get_google_embeddings() -> GoogleGenerativeAIEmbeddings`

* **Modelo**: `models/text-embedding-004`
* **Configuración**:
  * `transport='rest'`: REST más estable que gRPC
  * `batch_size=3`: Reducido para evitar rate limits
  * `output_dimensionality=768`: Dimensión de embeddings
  * `max_retries=3`: Reintentos limitados
  * `timeout=15`: Timeout aumentado

#### `get_cached_embeddings() -> GoogleGenerativeAIEmbeddings`

* **Patrón**: Singleton para evitar reinicializaciones
* **Cache**: Instancia global `_embeddings_instance`
* **Recovery**: Función `reset_embeddings_cache()`

#### `test_embeddings_health(embeddings) -> dict`

* **Propósito**: Verificar salud de embeddings
* **Métricas**: Latencia, dimensión, estado
* **Test**: Texto simple "servicio gubernamental República Dominicana"

### 5. Integración ChromaDB: `chroma_vector.py`

#### `get_chroma_ingest(docs, directory, collection_name, embeddings) -> VectorStoreRetriever`

* **Proceso**:
  1. Validar inputs
  2. Crear directorio de persistencia
  3. Limpiar metadata con `clean_metadata_for_chroma()`
  4. Crear vectorstore con `Chroma()`
  5. Ingestar en lotes de 50 documentos
  6. Crear retriever optimizado

#### `clean_metadata_for_chroma(docs) -> List[Document]`

* **Propósito**: Limpiar metadata para ChromaDB
* **Procesamiento**:
  * Strings, ints, floats, bools: Directo
  * Listas: Convertir a strings separados por comas
  * Dicts: Convertir a JSON string
  * Otros tipos: Convertir a string
* **Manejo de errores**: Continuar si falla un campo

#### `get_chroma_load(embeddings, directory, collection_name, retriever_k) -> VectorStoreRetriever`

* **Propósito**: Cargar ChromaDB existente
* **Validaciones**:
  * Directorio existe
  * Embeddings disponibles
  * Colección tiene documentos
* **Test**: Verificación básica con query "test"

#### `verify_chroma_health(directory, collection_name) -> dict`

* **Propósito**: Verificar salud de ChromaDB
* **Métricas**:
  * `directory_exists`: Directorio existe
  * `collection_exists`: Colección existe
  * `document_count`: Número de documentos
  * `status`: healthy/empty/error

### 6. Sistema de Logging: `logger.py`

#### `ProductionLogger`

* **Handlers**:
  * `errors.log`: Solo ERROR y CRITICAL (50MB, 5 backups)
  * `application.log`: INFO+ (100MB, 10 backups)
  * `debug.log`: DEBUG (solo en desarrollo, 200MB, 3 backups)
  * Console: ERROR en producción, INFO en desarrollo

#### Funciones de logging estructurado:

* `log_performance(operation, duration, **kwargs)`: Métricas de performance
* `log_user_interaction(user_query, response_length, processing_time)`: Interacciones de usuario
* `log_error_with_context(error, context)`: Errores con contexto
* `log_system_health(component, status, metrics)`: Salud del sistema

### 7. Rate Limiting: `rate_limiter.py`

#### `InMemoryRateLimiter`

* **Configuración**:
  * `requests_per_second=0.04`: 1 request cada 25 segundos
  * `check_every_n_seconds=0.1`: Verificar cada 100ms
  * `max_bucket_size=10`: Máximo burst size

#### `rate_limit(max_per_minute)`

* **Propósito**: Rate limiting personalizado
* **Cálculo**: `period = 60 / max_per_minute`
* **Comportamiento**: Sleep entre requests

## Dependencias entre módulos

### Flujo de dependencias

```
ingest.py
├── embeddings/gemini_embedding.py
│   ├── utils/logger.py
│   └── llm/utils.py
├── utils/doc_loader.py
│   └── utils/logger.py
├── utils/doc_split.py
│   └── utils/logger.py
└── vectors/chroma_vector.py
    ├── utils/logger.py
    └── langchain_chroma
```

### Imports críticos

```python
# ingest.py
from src.embeddings.gemini_embedding import get_google_embeddings
from src.utils.doc_loader import load_webpage, load_documents_from_directory, load_json_from_directory
from src.utils.doc_split import split_documents
from src.vectors.chroma_vector import get_chroma_ingest

# doc_loader.py
from langchain_community.document_loaders import PyPDFLoader, JSONLoader
from langchain_core.documents import Document
from langchain_unstructured import UnstructuredLoader
from src.utils.logger import logger

# chroma_vector.py
from langchain_chroma import Chroma
from langchain_core.documents import Document
from langchain_core.vectorstores import VectorStoreRetriever
from src.utils.logger import logger, log_error_with_context, log_performance, log_system_health
```

## Configuración y variables de entorno

### Variables requeridas

* `GOOGLE_API_KEY`: API key de Google para embeddings
* `ENV`: Entorno (production/development/staging)

### Configuración de ChromaDB

* **Directorio**: `storage/chroma/`
* **Colección**: `servicios` (por defecto)
* **Persistencia**: Automática con `persist_directory`

### Configuración de Embeddings

* **Modelo**: `text-embedding-004`
* **Dimensión**: 768
* **Transport**: REST
* **Batch size**: 3
* **Timeout**: 15 segundos

## Monitoreo y logging

### Logs estructurados

```json
{
  "timestamp": "2025-01-08T12:00:00",
  "level": "INFO",
  "module": "chroma_vector",
  "function": "get_chroma_ingest",
  "line": 158,
  "message": "PERFORMANCE: {\"type\": \"performance\", \"operation\": \"chroma_ingestion\", \"duration_ms\": 1868.28, \"docs_ingested\": 69, \"success\": true}",
  "pid": 1234
}
```

### Métricas de performance

* **metadata\_cleaning**: Tiempo de limpieza de metadata
* **chroma\_ingestion**: Tiempo total de ingesta
* **embeddings\_initialization**: Tiempo de inicialización de embeddings
* **chroma\_load**: Tiempo de carga de ChromaDB

### Health Checks

* **chromadb**: Estado de la base de datos vectorial
* **embeddings**: Estado de los embeddings
* **system**: Salud general del sistema

## Troubleshooting

### Problemas comunes

#### 1. Embeddings no inicializan

```python
# Verificar API key
import os
print(os.getenv("GOOGLE_API_KEY"))

# Test manual
from src.embeddings.gemini_embedding import get_google_embeddings
embeddings = get_google_embeddings()
```

#### 2. ChromaDB vacío

```python
# Verificar directorio
from pathlib import Path
chroma_dir = Path("storage/chroma")
print(chroma_dir.exists())
print(list(chroma_dir.iterdir()))

# Verificar colección
from src.vectors.chroma_vector import verify_chroma_health
health = verify_chroma_health(chroma_dir, "servicios")
print(health)
```

#### 3. Rate limiting

```python
# Verificar rate limiter
from src.utils.rate_limiter import rate_limiter
print(rate_limiter.requests_per_second)
```

### Comandos de diagnóstico

```bash
# Verificar logs
tail -f logs/application.log
tail -f logs/errors.log

# Verificar ChromaDB
ls -la storage/chroma/
du -sh storage/chroma/

# Test de embeddings
python3 -c "from src.embeddings.gemini_embedding import get_google_embeddings; print(get_google_embeddings())"
```

## Optimizaciones de performance

### 1. Carga en Lotes

* **ChromaDB**: Lotes de 50 documentos
* **Embeddings**: Batch size de 3
* **Rate limiting**: 0.04 requests/segundo

### 2. Cache de Embeddings

* **Singleton pattern**: Evita reinicializaciones
* **Cache global**: `_embeddings_instance`
* **Recovery**: Función de reset

### 3. Logging optimizado

* **Rotating files**: Evita archivos gigantes
* **Niveles diferenciados**: ERROR, INFO, DEBUG
* **JSON estructurado**: Fácil parsing

### 4. Manejo de errores

* **Graceful degradation**: Continúa si falla un lote
* **Context logging**: Errores con contexto
* **Health checks**: Verificación de estado


---

# Agent Instructions
This documentation is published with GitBook. GitBook is the documentation platform designed so that both humans and AI agents can read, navigate, and reason over technical content effectively. Learn more at gitbook.com.

## Querying This Documentation
If you need additional information that is not directly available in this page, you can query the documentation dynamically by asking a question.

Perform an HTTP GET request on the current page URL with the `ask` query parameter, and the optional `goal` query parameter:

```
GET https://public-intelligence.gitbook.io/taina-agente-ia-ogtic/arquitectura-y-conceptos/estructura-proyecto/sistema-ingesta.md?ask=<question>&goal=<endgoal>
```

`ask` is the immediate question: it should be specific, self-contained, and written in natural language.
`goal` is optional and describes the broader end goal you are ultimately trying to accomplish on behalf of the user. GitBook uses it to tailor the answer towards what is most useful for that goal.

The response will contain a direct answer to the question and relevant excerpts and sources from the documentation.

Use this mechanism when the answer is not explicitly present in the current page, you need clarification or additional context, or you want to retrieve related documentation sections.
