log_correlation
Implementa la correlación de logs entre los distintos microservicios del pipeline de verificación de identidad mediante identificadores compartidos como trace_id y session_id. Permite reconstruir el flujo completo de una sesión KYC desde la captura de selfie hasta la decisión final, facilitando el diagnóstico de fallos y la auditoría de verificaciones. Es fundamental para entender el comportamiento end-to-end del sistema distribuido.
When to use
Usar este skill cuando el observability_agent necesite implementar o mejorar la trazabilidad entre servicios del pipeline KYC, permitiendo seguir una sesión de verificación a través de liveness, OCR, face_match, antifraud y decision.
Instructions
Definir un middleware FastAPI que genere y propague el trace_id y session_id en cada request:
import uuid
from starlette.middleware.base import BaseHTTPMiddleware
from contextvars import ContextVar
trace_id_var: ContextVar[str] = ContextVar("trace_id", default="")
session_id_var: ContextVar[str] = ContextVar("session_id", default="")
class CorrelationMiddleware(BaseHTTPMiddleware):
async def dispatch(self, request, call_next):
trace_id = request.headers.get("X-Trace-ID", str(uuid.uuid4()))
session_id = request.headers.get("X-Session-ID", str(uuid.uuid4()))
trace_id_var.set(trace_id)
session_id_var.set(session_id)
response = await call_next(request)
response.headers["X-Trace-ID"] = trace_id
response.headers["X-Session-ID"] = session_id
return response
Configurar un logging filter que inyecte automáticamente los IDs de correlación en cada log:
import logging
class CorrelationFilter(logging.Filter):
def filter(self, record):
record.trace_id = trace_id_var.get("")
record.session_id = session_id_var.get("")
return True
logger = logging.getLogger("kyc-pipeline")
logger.addFilter(CorrelationFilter())
Configurar el formato de log JSON para incluir los campos de correlación:
import json
class JSONFormatter(logging.Formatter):
def format(self, record):
log_entry = {
"timestamp": self.formatTime(record),
"level": record.levelname,
"service": record.name,
"module": getattr(record, "module", "unknown"),
"trace_id": getattr(record, "trace_id", ""),
"session_id": getattr(record, "session_id", ""),
"message": record.getMessage(),
}
return json.dumps(log_entry)
Propagar los IDs de correlación en las llamadas HTTP entre microservicios:
import httpx
async def call_face_match(selfie_data: bytes, doc_face_data: bytes):
headers = {
"X-Trace-ID": trace_id_var.get(),
"X-Session-ID": session_id_var.get(),
}
async with httpx.AsyncClient() as client:
response = await client.post(
"http://face-match:8000/compare",
headers=headers,
files={"selfie": selfie_data, "document": doc_face_data}
)
return response.json()
Crear queries de correlación en Loki/LogQL para reconstruir una sesión completa:
{pipeline="kyc"} | json | session_id="sess-abc123" | line_format "{{.timestamp}} [{{.service}}] {{.message}}"
Crear queries equivalentes en Elasticsearch/Kibana:
{
"query": {
"bool": {
"must": [
{ "term": { "session_id": "sess-abc123" } }
]
}
},
"sort": [{ "timestamp": "asc" }]
}
Implementar un endpoint de diagnóstico que devuelva el timeline de una sesión:
@app.get("/api/v1/sessions/{session_id}/timeline")
async def get_session_timeline(session_id: str):
logs = await query_loki(f'{{pipeline="kyc"}} | json | session_id="{session_id}"')
return {
"session_id": session_id,
"steps": [parse_log_entry(log) for log in logs],
"total_duration_ms": calculate_duration(logs)
}
Notes
- El trace_id y session_id deben propagarse consistentemente en TODOS los microservicios del pipeline; un solo servicio sin propagación rompe la cadena de correlación.
- En entornos con colas de mensajes (Redis, RabbitMQ), los IDs de correlación deben incluirse en los metadatos del mensaje para mantener la trazabilidad asíncrona.
- Los IDs de correlación son esenciales para cumplir requisitos de auditoría GDPR/LOPD, ya que permiten reconstruir exactamente qué procesamiento se aplicó a los datos de un usuario.
1---2name: log-correlation3description: Correlación de logs entre microservicios del pipeline KYC usando trace_id y session_id.4---56# log_correlation78Implementa la correlación de logs entre los distintos microservicios del pipeline de verificación de identidad mediante identificadores compartidos como trace_id y session_id. Permite reconstruir el flujo completo de una sesión KYC desde la captura de selfie hasta la decisión final, facilitando el diagnóstico de fallos y la auditoría de verificaciones. Es fundamental para entender el comportamiento end-to-end del sistema distribuido.910## When to use1112Usar este skill cuando el observability_agent necesite implementar o mejorar la trazabilidad entre servicios del pipeline KYC, permitiendo seguir una sesión de verificación a través de liveness, OCR, face_match, antifraud y decision.1314## Instructions15161. Definir un middleware FastAPI que genere y propague el trace_id y session_id en cada request:17 ```python18 import uuid19 from starlette.middleware.base import BaseHTTPMiddleware20 from contextvars import ContextVar2122 trace_id_var: ContextVar[str] = ContextVar("trace_id", default="")23 session_id_var: ContextVar[str] = ContextVar("session_id", default="")2425 class CorrelationMiddleware(BaseHTTPMiddleware):26 async def dispatch(self, request, call_next):27 trace_id = request.headers.get("X-Trace-ID", str(uuid.uuid4()))28 session_id = request.headers.get("X-Session-ID", str(uuid.uuid4()))29 trace_id_var.set(trace_id)30 session_id_var.set(session_id)31 response = await call_next(request)32 response.headers["X-Trace-ID"] = trace_id33 response.headers["X-Session-ID"] = session_id34 return response35 ```36372. Configurar un logging filter que inyecte automáticamente los IDs de correlación en cada log:38 ```python39 import logging4041 class CorrelationFilter(logging.Filter):42 def filter(self, record):43 record.trace_id = trace_id_var.get("")44 record.session_id = session_id_var.get("")45 return True4647 logger = logging.getLogger("kyc-pipeline")48 logger.addFilter(CorrelationFilter())49 ```50513. Configurar el formato de log JSON para incluir los campos de correlación:52 ```python53 import json5455 class JSONFormatter(logging.Formatter):56 def format(self, record):57 log_entry = {58 "timestamp": self.formatTime(record),59 "level": record.levelname,60 "service": record.name,61 "module": getattr(record, "module", "unknown"),62 "trace_id": getattr(record, "trace_id", ""),63 "session_id": getattr(record, "session_id", ""),64 "message": record.getMessage(),65 }66 return json.dumps(log_entry)67 ```68694. Propagar los IDs de correlación en las llamadas HTTP entre microservicios:70 ```python71 import httpx7273 async def call_face_match(selfie_data: bytes, doc_face_data: bytes):74 headers = {75 "X-Trace-ID": trace_id_var.get(),76 "X-Session-ID": session_id_var.get(),77 }78 async with httpx.AsyncClient() as client:79 response = await client.post(80 "http://face-match:8000/compare",81 headers=headers,82 files={"selfie": selfie_data, "document": doc_face_data}83 )84 return response.json()85 ```86875. Crear queries de correlación en Loki/LogQL para reconstruir una sesión completa:88 ```logql89 {pipeline="kyc"} | json | session_id="sess-abc123" | line_format "{{.timestamp}} [{{.service}}] {{.message}}"90 ```91926. Crear queries equivalentes en Elasticsearch/Kibana:93 ```json94 {95 "query": {96 "bool": {97 "must": [98 { "term": { "session_id": "sess-abc123" } }99 ]100 }101 },102 "sort": [{ "timestamp": "asc" }]103 }104 ```1051067. Implementar un endpoint de diagnóstico que devuelva el timeline de una sesión:107 ```python108 @app.get("/api/v1/sessions/{session_id}/timeline")109 async def get_session_timeline(session_id: str):110 logs = await query_loki(f'{{pipeline="kyc"}} | json | session_id="{session_id}"')111 return {112 "session_id": session_id,113 "steps": [parse_log_entry(log) for log in logs],114 "total_duration_ms": calculate_duration(logs)115 }116 ```117118## Notes119120- El trace_id y session_id deben propagarse consistentemente en TODOS los microservicios del pipeline; un solo servicio sin propagación rompe la cadena de correlación.121- En entornos con colas de mensajes (Redis, RabbitMQ), los IDs de correlación deben incluirse en los metadatos del mensaje para mantener la trazabilidad asíncrona.122- Los IDs de correlación son esenciales para cumplir requisitos de auditoría GDPR/LOPD, ya que permiten reconstruir exactamente qué procesamiento se aplicó a los datos de un usuario.