Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2f990ea636 | ||
|
|
4244b00ce7 | ||
|
|
dd5cfd9799 | ||
|
|
9377730de4 | ||
|
|
ce629139c2 | ||
|
|
57ba136c70 | ||
|
|
9dd2a8f9bd | ||
|
|
c860984ae4 | ||
|
|
c414e16357 | ||
|
|
eac2ed02a4 |
@@ -28,7 +28,16 @@ necesita antes de operar agentes con impacto real.
|
||||
└────────────────────────────────────────────────────────────────────┘
|
||||
```
|
||||
|
||||
Detalles completos en [`ARCHITECTURE.md`](ARCHITECTURE.md).
|
||||
Detalles completos en [`ARCHITECTURE.md`](ARCHITECTURE.md). Además:
|
||||
- [`docs/walkthrough.html`](docs/walkthrough.html) — **walkthrough HTML autocontenido**
|
||||
(de alto a bajo nivel, con diagramas). Ábrelo directamente en el navegador
|
||||
(`xdg-open docs/walkthrough.html`); no requiere servidor.
|
||||
- [`docs/explicacion.md`](docs/explicacion.md) — explicación didáctica de extremo a
|
||||
extremo (conceptos, recorrido por todos los módulos y sus interrelaciones, y el
|
||||
viaje de una petición de principio a fin).
|
||||
- [`docs/componentes.md`](docs/componentes.md) — referencia de cableado a bajo nivel
|
||||
(grafo de dependencias de módulos, inyección de dependencias, firmas de los
|
||||
contratos entre capas, cadenas de llamada de cada endpoint).
|
||||
|
||||
## Quickstart
|
||||
|
||||
|
||||
@@ -147,8 +147,47 @@ async def invoke_agent(
|
||||
|
||||
|
||||
@router.get("", response_model=list[AgentExecutionSummary])
|
||||
def list_executions(settings: SettingsDep) -> list[AgentExecutionSummary]:
|
||||
return read_execution_summaries(settings.data_dir)
|
||||
async def list_executions(
|
||||
registry: RegistryDep,
|
||||
policies: PolicyStoreDep,
|
||||
orchestrator: OrchestratorDep,
|
||||
settings: SettingsDep,
|
||||
) -> list[AgentExecutionSummary]:
|
||||
"""Resumen de ejecuciones: las terminales del JSONL más las pausadas en HITL.
|
||||
|
||||
Una ejecución en ``awaiting_approval`` no se escribe en el log append-only
|
||||
(`executions.jsonl`); solo vive en el checkpointer. Para que la página de
|
||||
aprobaciones del dashboard la encuentre, aquí se reconstruyen esas desde el
|
||||
índice `execution_index.json` + el checkpointer.
|
||||
"""
|
||||
summaries = read_execution_summaries(settings.data_dir)
|
||||
seen = {str(s.trace_id) for s in summaries}
|
||||
for trace_id, meta in _load_index(settings.data_dir).items():
|
||||
if trace_id in seen:
|
||||
continue
|
||||
try:
|
||||
agent_def = registry.get_version(meta["agent_name"], meta["version"])
|
||||
policy = policies.get_policy(agent_def.guardrails[0])
|
||||
execution = await orchestrator.snapshot(
|
||||
agent_def=agent_def, policy=policy, trace_id=UUID(trace_id)
|
||||
)
|
||||
except (FileNotFoundError, IndexError, KeyError, ValueError):
|
||||
continue
|
||||
if execution is None:
|
||||
continue
|
||||
summaries.append(
|
||||
AgentExecutionSummary(
|
||||
trace_id=execution.trace_id,
|
||||
agent_name=execution.agent_name,
|
||||
agent_version=execution.agent_version,
|
||||
status=execution.status,
|
||||
started_at=execution.started_at,
|
||||
finished_at=execution.finished_at,
|
||||
n_violations=len(execution.violations),
|
||||
n_proposed_actions=len(execution.proposed_actions),
|
||||
)
|
||||
)
|
||||
return summaries
|
||||
|
||||
|
||||
@router.get("/{trace_id}", response_model=AgentExecution)
|
||||
|
||||
@@ -21,6 +21,33 @@ _INJECTION_PATTERNS = [
|
||||
r"reveal (the )?(system|hidden) (prompt|instruction)",
|
||||
]
|
||||
|
||||
# Singleton perezoso del AnalyzerEngine de Presidio. Construirlo es caro (carga el
|
||||
# modelo spaCy y los recognizers), así que se reutiliza entre llamadas.
|
||||
_PRESIDIO_ANALYZER: Any = None
|
||||
|
||||
|
||||
def _presidio_analyzer() -> Any:
|
||||
"""``AnalyzerEngine`` de Presidio configurado con el modelo spaCy ``en_core_web_sm``.
|
||||
|
||||
Presidio usa por defecto ``en_core_web_lg`` (~560 MB), que **no** está en la imagen
|
||||
Docker: ``core/Dockerfile`` instala ``en_core_web_sm``. Si Presidio no está instalado
|
||||
o el modelo no se puede cargar, esto lanza y el llamador (`detect_pii`) hace fallback
|
||||
a la detección por regex.
|
||||
"""
|
||||
global _PRESIDIO_ANALYZER
|
||||
if _PRESIDIO_ANALYZER is None:
|
||||
from presidio_analyzer import AnalyzerEngine
|
||||
from presidio_analyzer.nlp_engine import NlpEngineProvider
|
||||
|
||||
nlp_engine = NlpEngineProvider(
|
||||
nlp_configuration={
|
||||
"nlp_engine_name": "spacy",
|
||||
"models": [{"lang_code": "en", "model_name": "en_core_web_sm"}],
|
||||
}
|
||||
).create_engine()
|
||||
_PRESIDIO_ANALYZER = AnalyzerEngine(nlp_engine=nlp_engine)
|
||||
return _PRESIDIO_ANALYZER
|
||||
|
||||
|
||||
def _violation(
|
||||
*,
|
||||
@@ -45,11 +72,11 @@ def _violation(
|
||||
def detect_pii(
|
||||
text: str, config: dict[str, Any], trace_id: UUID, stage: str
|
||||
) -> list[GuardrailViolation]:
|
||||
"""Detección PII vía Presidio Analyzer (con fallback a regex si no está instalado)."""
|
||||
"""Detección PII vía Presidio Analyzer (con fallback a regex si no está disponible)."""
|
||||
try:
|
||||
from presidio_analyzer import AnalyzerEngine
|
||||
analyzer = _presidio_analyzer()
|
||||
except Exception:
|
||||
# Si Presidio no está, fallback a regex básica
|
||||
# Presidio no instalado o modelo spaCy no disponible → fallback a regex básica.
|
||||
return _pii_regex_fallback(text, config, trace_id, stage)
|
||||
|
||||
entities = config.get(
|
||||
@@ -59,7 +86,6 @@ def detect_pii(
|
||||
sev = config.get("severity_on_match", "block")
|
||||
blocked = sev == "block"
|
||||
|
||||
analyzer = AnalyzerEngine()
|
||||
results = analyzer.analyze(text=text, entities=entities, language="en")
|
||||
if not results:
|
||||
return []
|
||||
|
||||
@@ -66,6 +66,9 @@ class CoreClient:
|
||||
def _post(self, path: str, body: dict) -> Any:
|
||||
r = self._client.post(path, json=body)
|
||||
if r.status_code in {404, 409, 422}:
|
||||
return {"error": r.json()}
|
||||
# Clave deliberadamente distinta de "error": el cuerpo de un AgentExecution
|
||||
# correcto ya trae error=None, así que las páginas no podrían distinguir
|
||||
# "ejecución sin error" de "la API rechazó la petición" si reusáramos "error".
|
||||
return {"api_error": r.json()}
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
|
||||
@@ -53,7 +53,9 @@ if st.button("🚀 Invocar agente", type="primary", disabled=not user_input):
|
||||
st.session_state["last_execution"] = result
|
||||
|
||||
execution = st.session_state.get("last_execution")
|
||||
if execution and "error" not in execution:
|
||||
if execution and "api_error" in execution:
|
||||
st.error(f"Error de la API: {execution['api_error']}")
|
||||
elif execution:
|
||||
st.divider()
|
||||
st.markdown(f"**trace_id:** `{execution['trace_id']}`")
|
||||
st.markdown(f"**Status:** `{execution['status']}`")
|
||||
@@ -63,10 +65,10 @@ if execution and "error" not in execution:
|
||||
"Esta ejecución requiere aprobación humana. "
|
||||
"Ve a la página **Aprobaciones** para revisar y decidir."
|
||||
)
|
||||
if execution.get("error"):
|
||||
st.error(f"La ejecución terminó con error: `{execution['error']}`")
|
||||
if execution.get("final_output"):
|
||||
st.subheader("📦 Output final")
|
||||
st.json(execution["final_output"])
|
||||
render_violations(execution.get("violations", []))
|
||||
render_trace(execution.get("decision_path", []))
|
||||
elif execution and "error" in execution:
|
||||
st.error(f"Error de la API: {execution['error']}")
|
||||
|
||||
@@ -65,8 +65,8 @@ with col_a:
|
||||
execution["trace_id"],
|
||||
{"approved_action_ids": approved_ids, "comment": comment},
|
||||
)
|
||||
if "error" in r:
|
||||
st.error(r["error"])
|
||||
if "api_error" in r:
|
||||
st.error(r["api_error"])
|
||||
else:
|
||||
st.success(f"Status: {r['status']}")
|
||||
st.rerun()
|
||||
@@ -75,8 +75,8 @@ with col_r:
|
||||
if st.button("❌ Rechazar ejecución", disabled=not reason):
|
||||
with st.spinner("Aplicando rechazo..."):
|
||||
r = client.reject(execution["trace_id"], {"reason": reason})
|
||||
if "error" in r:
|
||||
st.error(r["error"])
|
||||
if "api_error" in r:
|
||||
st.error(r["api_error"])
|
||||
else:
|
||||
st.success(f"Status: {r['status']}")
|
||||
st.rerun()
|
||||
|
||||
@@ -0,0 +1,427 @@
|
||||
# Componentes de AgentForge y su interrelación (bajo nivel)
|
||||
|
||||
> **Alcance.** Este documento es la **referencia de cableado**: módulos exactos,
|
||||
> firmas, el grafo de dependencias de imports, el grafo de inyección de
|
||||
> dependencias, los contratos entre capas y las cadenas de llamada de cada
|
||||
> endpoint. Es preciso, no narrativo.
|
||||
>
|
||||
> - ¿Quieres la historia y el "por qué"? → [`docs/explicacion.md`](explicacion.md).
|
||||
> - ¿Las decisiones técnicas resumidas? → [`ARCHITECTURE.md`](../ARCHITECTURE.md).
|
||||
> - ¿Cómo arrancarlo? → [`README.md`](../README.md).
|
||||
>
|
||||
> Rutas relativas a `core/src/agentforge_core/` salvo que se diga otra cosa.
|
||||
> Refleja el estado del repo en `v0.1.0`.
|
||||
|
||||
---
|
||||
|
||||
## 1. Grafo de dependencias de módulos (imports internos)
|
||||
|
||||
Edges = "X importa de Y". El grafo es un DAG; las capas de abajo no importan nada
|
||||
de las de arriba. Nivel = profundidad topológica.
|
||||
|
||||
```
|
||||
NIVEL 0 (no importan nada del proyecto)
|
||||
config ← Settings (pydantic-settings, lee .env)
|
||||
observability.logging ← configure_logging / bind_trace_id / clear_trace_id
|
||||
domain.agent ← AgentDefinition, LLMConfig, AgentVersionMeta
|
||||
domain.guardrail ← GuardrailViolation
|
||||
domain.policy ← PolicyDefinition, PolicyValidator, PolicyVersionMeta
|
||||
registry.versioning ← compute_hash, unified_diff, DiffResult
|
||||
llm.base ← LLMProvider (Protocol), Message, CompletionResult
|
||||
runtime.state ← AgentState (TypedDict)
|
||||
runtime.checkpointer ← build_checkpointer (AsyncSqliteSaver)
|
||||
|
||||
NIVEL 1
|
||||
domain.execution ⇐ domain.guardrail
|
||||
llm.mock / llm.azure / llm.openai ⇐ llm.base
|
||||
registry.repository ⇐ domain.agent, registry.versioning
|
||||
registry.policy_store ⇐ domain.policy
|
||||
guardrails.validators ⇐ domain.guardrail
|
||||
guardrails.base ⇐ domain.guardrail, domain.policy
|
||||
api.persistence ⇐ domain.execution, domain.guardrail
|
||||
api.middlewares ⇐ observability.logging
|
||||
|
||||
NIVEL 2
|
||||
llm.factory ⇐ config, llm.base, llm.{mock,azure,openai}
|
||||
registry.factory ⇐ config, registry.{repository,policy_store}
|
||||
guardrails.guardrails_ai ⇐ domain.{guardrail,policy}, guardrails.validators
|
||||
guardrails.composite ⇐ domain.{guardrail,policy}, guardrails.base
|
||||
guardrails.nemo ⇐ domain.{guardrail,policy}
|
||||
|
||||
NIVEL 3
|
||||
guardrails.factory ⇐ config, guardrails.{base,composite,guardrails_ai,nemo}
|
||||
runtime.nodes ⇐ domain.{agent,policy}, guardrails.base, llm.base, runtime.state
|
||||
runtime.graph ⇐ domain.{agent,policy}, guardrails.base, llm.base, runtime.{nodes,state}
|
||||
|
||||
NIVEL 4
|
||||
runtime.orchestrator ⇐ domain.{agent,execution,guardrail,policy}, guardrails.base, llm.base, runtime.{checkpointer,graph}
|
||||
|
||||
NIVEL 5
|
||||
api.deps ⇐ config, guardrails.{base,factory}, llm.{base,factory}, registry.{factory,policy_store,repository}, runtime.orchestrator
|
||||
api.agents ⇐ api.deps, domain.agent, registry.versioning
|
||||
api.policies ⇐ api.deps, domain.policy
|
||||
api.violations ⇐ api.deps, api.persistence, domain.guardrail
|
||||
api.executions ⇐ api.deps, api.persistence, domain.{agent,execution,policy}, registry.{policy_store,repository}, runtime.orchestrator
|
||||
|
||||
NIVEL 6
|
||||
main ⇐ api (routers), api.middlewares, config, observability.logging → create_app(), app
|
||||
|
||||
(separado, sin imports del paquete core)
|
||||
dashboard/src/agentforge_dashboard/* ← habla con `main` por HTTP, no por import
|
||||
```
|
||||
|
||||
Reglas que se cumplen y conviene mantener:
|
||||
- **`domain/` no importa nada del resto del proyecto.** Es el vocabulario; todos dependen de él.
|
||||
- **Solo `api/` importa de `runtime.orchestrator`** (y `deps.py` lo construye). El resto de `api/` no toca LangGraph; `runtime/` no toca `api/`.
|
||||
- **Solo los `factory.py` y `main.py` importan `config.Settings`.** El resto recibe los objetos ya construidos.
|
||||
- **`guardrails.composite` no conoce a sus sub-engines concretos** (solo `GuardrailEngine`); los ensambla `guardrails.factory`.
|
||||
|
||||
---
|
||||
|
||||
## 2. Inyección de dependencias — `api/deps.py`
|
||||
|
||||
Cinco objetos del dominio, cada uno construido **una sola vez** por proceso
|
||||
(`@lru_cache(maxsize=1)`), más el orchestrator que los compone:
|
||||
|
||||
```
|
||||
get_settings() ──────────────────────────────────────────────► Settings() (pydantic-settings ← .env)
|
||||
│
|
||||
├──► get_registry() = build_agent_registry(settings) ──► FileSystemAgentRegistry(settings.agents_dir)
|
||||
├──► get_policy_store() = build_policy_store(settings) ──► FileSystemPolicyStore(settings.policies_dir)
|
||||
├──► get_llm_provider() = build_llm_provider(settings) ──► Mock | AzureOpenAI | OpenAI (según settings.llm_provider)
|
||||
├──► get_guardrail_engine() = build_guardrail_engine(settings)──► CompositeGuardrailEngine([GuardrailsAIEngine(), (NeMoGuardrailsEngine() si settings.guardrails_nemo_enabled)])
|
||||
│
|
||||
└──► get_orchestrator() = AgentOrchestrator(
|
||||
provider = get_llm_provider(),
|
||||
engine = get_guardrail_engine(),
|
||||
data_dir = get_settings().data_dir,
|
||||
)
|
||||
```
|
||||
|
||||
Aliases que los routers piden por parámetro (`Annotated[T, Depends(get_*)]`):
|
||||
|
||||
| Alias | Tipo | Lo usan |
|
||||
|-------|------|---------|
|
||||
| `SettingsDep` | `Settings` | `executions` (todos los endpoints), `violations` |
|
||||
| `RegistryDep` | `FileSystemAgentRegistry` | `agents` (todos), `executions` (`invoke`, `get`, `list`, `approve`, `reject`) |
|
||||
| `PolicyStoreDep` | `FileSystemPolicyStore` | `policies` (todos), `executions` (`invoke`, `get`, `list`, `approve`, `reject`) |
|
||||
| `OrchestratorDep` | `AgentOrchestrator` | `executions` (`invoke`, `get`, `list`, `approve`, `reject`) |
|
||||
|
||||
> No hay `LLMProviderDep` ni `GuardrailEngineDep`: el provider y el engine **solo**
|
||||
> se inyectan al `AgentOrchestrator` (vía `get_orchestrator`), nunca a un router.
|
||||
>
|
||||
> Lifecycle: las `get_*` son perezosas → se llaman en la primera request que las
|
||||
> necesita. Los tests hacen `deps.get_settings.cache_clear()` (etc.) tras
|
||||
> `monkeypatch.setenv("DATA_DIR", tmp_path)` para reconstruir todo apuntando a un
|
||||
> directorio temporal. `main.create_app()` llama a `Settings()` *directamente*
|
||||
> (para configurar el logging), no a `deps.get_settings()`; son instancias
|
||||
> distintas pero leen la misma config.
|
||||
|
||||
---
|
||||
|
||||
## 3. Contratos entre componentes (firmas exactas)
|
||||
|
||||
### 3.1 Proveedor LLM — `llm/base.py`
|
||||
|
||||
```python
|
||||
class Message(BaseModel): role: Literal["system","user","assistant"]; content: str
|
||||
class CompletionResult(BaseModel): content: str; model: str; tokens_in: int; tokens_out: int; latency_ms: int
|
||||
|
||||
class LLMProvider(Protocol):
|
||||
name: str
|
||||
async def complete(self, messages: list[Message], schema: dict[str,Any] | None = None,
|
||||
temperature: float = 0.2, max_tokens: int = 2000) -> CompletionResult: ...
|
||||
```
|
||||
|
||||
Implementaciones: `MockProvider` (`llm/mock.py`, determinista — elige una de
|
||||
`_CANONICAL_RESPONSES = {"sip":…, "mos":…, "hss":…}` buscando esas subcadenas en el
|
||||
input, en ese orden; si no, una respuesta genérica), `AzureOpenAIProvider`
|
||||
(`llm/azure.py`), `OpenAIProvider` (`llm/openai.py`). Selector: `build_llm_provider(settings)`
|
||||
en `llm/factory.py` (`match settings.llm_provider`). El nodo `llm_reason` lo invoca
|
||||
así: `await provider.complete(messages=[Message("system", agent_def.system_prompt), Message("user", state["user_input"])], temperature=agent_def.llm.temperature, max_tokens=agent_def.llm.max_tokens)`.
|
||||
|
||||
### 3.2 Motor de guardrails — `guardrails/base.py`
|
||||
|
||||
```python
|
||||
class GuardrailEngine(Protocol):
|
||||
name: str
|
||||
async def validate_input (self, payload: str, policy: PolicyDefinition, trace_id: UUID) -> list[GuardrailViolation]: ...
|
||||
async def validate_output(self, payload: dict[str, Any], policy: PolicyDefinition, trace_id: UUID) -> list[GuardrailViolation]: ...
|
||||
```
|
||||
|
||||
```python
|
||||
class GuardrailViolation(BaseModel):
|
||||
trace_id: UUID; timestamp: datetime
|
||||
stage: Literal["input","output"]; validator: str
|
||||
severity: Literal["info","warning","block"]; message: str; blocked: bool
|
||||
```
|
||||
|
||||
Implementaciones:
|
||||
- **`CompositeGuardrailEngine(engines: list[GuardrailEngine])`** (`guardrails/composite.py`) — `validate_*` hace `asyncio.gather(*(e.validate_*(...) for e in engines))` y aplana las listas. Lanza `ValueError` si `engines` está vacío.
|
||||
- **`GuardrailsAIEngine()`** (`guardrails/guardrails_ai.py`) — el real. Mantiene dos registros `type → función`:
|
||||
- `INPUT_VALIDATORS = {"detect_pii", "prompt_injection", "toxic_language", "forbidden_topics"}`
|
||||
- `OUTPUT_VALIDATORS = {"schema_match", "pii_leakage", "forbidden_action_keywords", "telco_safety_rules"}`
|
||||
|
||||
`_run(...)` recorre `policy.input_validators` / `policy.output_validators` (cada uno un `PolicyValidator(type, config)`), busca `type` en el registro, y llama `fn(payload, validator.config, trace_id, kind)`. Si `fn` lanza y `policy.on_validator_error == "fail_closed"` ⇒ añade una `GuardrailViolation(severity="block", blocked=True, message=f"validator failed: {exc}", validator=v.type)`. Si el `type` no existe en el registro ⇒ log warning y continúa (no bloquea).
|
||||
- **`NeMoGuardrailsEngine(allowed_keywords)`** (`guardrails/nemo.py`) — stub; solo se añade al composite si `settings.guardrails_nemo_enabled`. `validate_input`: si el texto no menciona ningún keyword permitido ⇒ una violación `severity="warning"`, `blocked=False`. `validate_output`: siempre `[]`.
|
||||
|
||||
Selector: `build_guardrail_engine(settings)` en `guardrails/factory.py`.
|
||||
|
||||
**Contrato de una función validadora** (todas las de `guardrails/validators.py`):
|
||||
|
||||
```python
|
||||
def <validador>(payload: <str|dict>, config: dict[str, Any], trace_id: UUID, stage: str) -> list[GuardrailViolation]
|
||||
```
|
||||
|
||||
- Entrada (`validate_input`): `payload` es `str` (el texto del usuario). `detect_pii`, `prompt_injection`, `toxic_language`, `forbidden_topics`.
|
||||
- Salida (`validate_output`): `payload` es `dict` (el JSON parseado del LLM). `schema_match`, `pii_leakage`, `forbidden_action_keywords`, `telco_safety_rules`.
|
||||
- `config` viene literal del YAML de la política (p. ej. `{entities:[...], severity_on_match:"block"}`).
|
||||
- `detect_pii`: usa `_presidio_analyzer()` (singleton perezoso de `presidio_analyzer.AnalyzerEngine` con el modelo spaCy `en_core_web_sm`); si Presidio no está instalado o el modelo no carga, hace fallback a `_pii_regex_fallback` (regex de EMAIL / PHONE 3-3-3 / ES_NIF `\d{8}[A-HJ-NP-TV-Z]` / IP).
|
||||
|
||||
### 3.3 Registry de agentes — `registry/repository.py` → `FileSystemAgentRegistry(root: Path)`
|
||||
|
||||
| Método | Devuelve | Lee |
|
||||
|--------|----------|-----|
|
||||
| `list_agents()` | `list[AgentDefinition]` | recorre `root/*/index.yaml`, devuelve la versión activa de cada uno |
|
||||
| `get_agent(name, version=None)` | `AgentDefinition` | `index.yaml["active_version"]` (o `version`) → `get_version` |
|
||||
| `get_version(name, version)` | `AgentDefinition` | `root/<name>/versions/<version>.yaml` (lanza `FileNotFoundError`) |
|
||||
| `list_versions(name)` | `list[AgentVersionMeta]` | `index.yaml["versions"]` |
|
||||
| `upsert_version(name, body, message, author)` | `AgentVersionMeta` | escribe `versions/<id>.yaml` (con `compute_hash`) y actualiza `index.yaml` (si `body.state=="active"` cambia `active_version`) |
|
||||
| `diff_versions(name, v1, v2)` | `DiffResult` | `versioning.unified_diff` entre los dos YAML |
|
||||
|
||||
### 3.4 Store de políticas — `registry/policy_store.py` → `FileSystemPolicyStore(root: Path)`
|
||||
|
||||
`list_policies() -> list[PolicyDefinition]`, `get_policy(name, version=None) -> PolicyDefinition` (lanza `FileNotFoundError`), `list_versions(name) -> list[PolicyVersionMeta]`.
|
||||
|
||||
```python
|
||||
class PolicyValidator(BaseModel): type: str; config: dict[str,Any] = {}
|
||||
class PolicyDefinition(BaseModel):
|
||||
name: str; version: str; description: str
|
||||
input_validators: list[PolicyValidator]
|
||||
output_validators: list[PolicyValidator]
|
||||
on_validator_error: Literal["fail_open","fail_closed"] = "fail_closed"
|
||||
class PolicyVersionMeta(BaseModel): id: str; hash: str; author: str; message: str; created_at: datetime
|
||||
```
|
||||
|
||||
### 3.5 Orchestrator — `runtime/orchestrator.py` → `AgentOrchestrator`
|
||||
|
||||
```python
|
||||
AgentOrchestrator(*, provider: LLMProvider, engine: GuardrailEngine, data_dir: Path)
|
||||
|
||||
async invoke (*, agent_def: AgentDefinition, policy: PolicyDefinition, user_input: str,
|
||||
trace_id: UUID | None = None) -> AgentExecution
|
||||
async resume (*, agent_def, policy, trace_id: UUID, decision: dict[str,Any]) -> AgentExecution
|
||||
async snapshot(*, agent_def, policy, trace_id: UUID) -> AgentExecution | None
|
||||
```
|
||||
|
||||
- Cada uno abre su **propio** `AsyncSqliteSaver` (`runtime/checkpointer.build_checkpointer(data_dir)` → `data_dir/checkpoints.sqlite`, context manager async) y construye el grafo con `build_graph(...)`.
|
||||
- `invoke`: arma el `AgentState` inicial (`status="running"`, listas vacías…) y hace `graph.ainvoke(state, config={"configurable": {"thread_id": str(trace_id)}})`. Si el grafo lanza, marca `crashed=True`. Luego `_snapshot(...)` → `AgentExecution`.
|
||||
- `resume`: `graph.ainvoke(Command(resume=decision), config={thread_id: trace_id})` — reanuda el `interrupt()` con `decision`.
|
||||
- `snapshot`: `graph.aget_state(config)`; si `state.values` está vacío → `None`; si no, `_build_execution(...)`.
|
||||
|
||||
### 3.6 El grafo — `runtime/graph.py` + `runtime/nodes.py` + `runtime/state.py`
|
||||
|
||||
```python
|
||||
# runtime/state.py
|
||||
class AgentState(TypedDict, total=False):
|
||||
trace_id: str; agent_name: str; agent_version: str; user_input: str
|
||||
messages: list[dict]; raw_llm_output: str | None; parsed_output: dict | None
|
||||
proposed_actions: list[dict]; violations: list[dict]
|
||||
decision_path: Annotated[list[dict], operator.add] # ← reducer: cada nodo AÑADE pasos
|
||||
status: str; error: str | None
|
||||
human_decision: dict | None; final_output: dict | None
|
||||
```
|
||||
|
||||
```python
|
||||
# runtime/nodes.py
|
||||
NodeFn = Callable[[AgentState], Awaitable[dict[str, Any]]]
|
||||
# factories que devuelven un NodeFn, parametrizadas con engine/policy/provider/agent_def:
|
||||
build_node_validate_input(engine, policy) build_node_llm_reason(provider, agent_def)
|
||||
build_node_validate_output(engine, policy) build_node_propose_actions()
|
||||
build_node_approve_gate(agent_def) build_node_finalize()
|
||||
# cada nodo añade un DecisionStep a decision_path: {step, timestamp, duration_ms, detail}
|
||||
```
|
||||
|
||||
```python
|
||||
# runtime/graph.py
|
||||
build_graph(*, agent_def, policy, provider, engine, checkpointer: BaseCheckpointSaver) -> CompiledGraph
|
||||
```
|
||||
|
||||
Cableado (nodos y aristas):
|
||||
|
||||
```
|
||||
START ─► validate_input
|
||||
validate_input ──(conditional)──► END si status == "blocked_by_guardrail"
|
||||
llm_reason en otro caso
|
||||
llm_reason ──(conditional)──────► END si status == "failed"
|
||||
validate_output en otro caso
|
||||
validate_output ──(conditional)─► END si status in {"blocked_by_guardrail","failed"}
|
||||
propose_actions en otro caso
|
||||
propose_actions ──► approve_gate ──► finalize ──► END (aristas fijas)
|
||||
```
|
||||
|
||||
- `approve_gate`: si hay acciones con `risk_score >= agent_def.risk_threshold_for_hitl` **o** `requires_approval` ⇒ `decision = interrupt({"awaiting_actions": risky})` → el grafo se **pausa** ahí (LangGraph persiste el estado en el checkpointer); al reanudar con `Command(resume=decision)`, `interrupt()` devuelve `decision` y el nodo continúa, escribiendo `human_decision`.
|
||||
- `finalize`: si `human_decision.rejected` ⇒ `status="failed"`, `error="rejected_by_human"`. Si no, `final_output = {**parsed_output, "approved_actions": <acciones cuyo id está en approved_action_ids; o todas si no hubo HITL>}`, `status="completed"`.
|
||||
|
||||
### 3.7 De `StateSnapshot` a `AgentExecution` — `orchestrator._build_execution`
|
||||
|
||||
```python
|
||||
class AgentExecution(BaseModel):
|
||||
trace_id: UUID; agent_name: str; agent_version: str
|
||||
status: Literal["running","awaiting_approval","blocked_by_guardrail","completed","failed"]
|
||||
started_at: datetime; finished_at: datetime | None
|
||||
decision_path: list[DecisionStep]; violations: list[GuardrailViolation]
|
||||
proposed_actions: list[ProposedAction]; needs_human_for: list[ProposedAction] | None
|
||||
final_output: dict[str,Any] | None; error: str | None
|
||||
class DecisionStep(BaseModel): step: str; timestamp: datetime; duration_ms: int; detail: dict[str,Any]
|
||||
class ProposedAction(BaseModel): id: str; action: str; target: str; risk_score: int (1..5); rollback_plan: str; requires_approval: bool
|
||||
class AgentExecutionSummary(BaseModel): trace_id; agent_name; agent_version; status; started_at; finished_at; n_violations: int; n_proposed_actions: int
|
||||
```
|
||||
|
||||
Reglas para derivar `status` a partir del `StateSnapshot` de LangGraph:
|
||||
1. `status = state.values.get("status", "running")`.
|
||||
2. Si `crashed` y `status` no es terminal ⇒ `status = "failed"` (`error = error_del_state or "internal_error"`).
|
||||
3. Si `state.next` (hay un nodo pendiente) y `status` no es terminal ⇒ `status = "awaiting_approval"` (es el `interrupt()` de `approve_gate`).
|
||||
4. `needs_human_for`: solo si `status == "awaiting_approval"` ⇒ `[a for a in proposed_actions if a.risk_score >= agent_def.risk_threshold_for_hitl or a.requires_approval]`; en otro caso `None`.
|
||||
5. `started_at = decision_path[0].timestamp` (o `now` si vacío); `finished_at = now` solo si `status` es terminal.
|
||||
6. Terminales = `{"completed","failed","blocked_by_guardrail"}`.
|
||||
|
||||
---
|
||||
|
||||
## 4. Cadenas de llamada por endpoint (`api/`)
|
||||
|
||||
Todos pasan primero por `TraceIdMiddleware` (lee/crea `X-Trace-Id`, lo bind-ea a structlog, lo devuelve en la respuesta). Montaje (en `main.create_app`): `agents.router→/agents`, `executions.invoke_router→/agents`, `executions.router→/executions`, `policies.router→/policies`, `violations.router→/violations`, más `GET /health`.
|
||||
|
||||
| Endpoint | Handler (`api/…`) | Deps inyectadas | Llama a → Escribe |
|
||||
|----------|-------------------|-----------------|-------------------|
|
||||
| `POST /agents/{name}/invoke` | `executions.invoke_agent` | Registry, PolicyStore, Orchestrator, Settings | `registry.get_agent(name)` → `policies.get_policy(agent.guardrails[0])` → `orchestrator.invoke(agent_def, policy, body.input)` → **`_record_execution(data_dir, trace_id, name, version)`** (escribe `execution_index.json`) → si `status` terminal: **`append_execution`** (`executions.jsonl`) → por cada violación: **`append_violation`** (`violations.jsonl`). 404 si el agente no existe; 422 si no tiene política. |
|
||||
| `GET /executions` | `executions.list_executions` | Registry, PolicyStore, Orchestrator, Settings | `read_execution_summaries(data_dir)` (lee `executions.jsonl`) + por cada `trace_id` en `execution_index.json` no presente ya: `registry.get_version` + `policies.get_policy` + `orchestrator.snapshot` → `AgentExecutionSummary`; fusiona. (Así aparecen también las `awaiting_approval`, que no están en el JSONL.) |
|
||||
| `GET /executions/{trace_id}` | `executions.get_execution` | Registry, PolicyStore, Orchestrator, Settings | `_resolve(...)` (lee `execution_index.json` → reconstruye `agent_def`/`policy`; 404 si no está, 500 si la config referida no existe) → `orchestrator.snapshot(...)` → 404 si `None`. |
|
||||
| `POST /executions/{trace_id}/approve` | `executions.approve_execution` | Registry, PolicyStore, Orchestrator, Settings | `_resolve` → `_ensure_awaiting` (`orchestrator.snapshot`; 409 si no está en `awaiting_approval`) → `orchestrator.resume(decision={approved_action_ids: body.approved_action_ids, comment: body.comment, rejected: False})` → **`append_execution`** (`executions.jsonl`). |
|
||||
| `POST /executions/{trace_id}/reject` | `executions.reject_execution` | Registry, PolicyStore, Orchestrator, Settings | `_resolve` → `_ensure_awaiting` → `orchestrator.resume(decision={approved_action_ids: [], rejected: True, reason: body.reason})` → **`append_execution`**. |
|
||||
| `GET /agents` | `agents.list_agents` | Registry | `registry.list_agents()` |
|
||||
| `GET /agents/{name}` | `agents.get_agent` | Registry | `registry.get_agent(name)` (404 si no existe) |
|
||||
| `GET /agents/{name}/versions` | `agents.list_versions` | Registry | `registry.list_versions(name)` (404) |
|
||||
| `GET /agents/{name}/versions/{version}` | `agents.get_version` | Registry | `registry.get_version(name, version)` (404) |
|
||||
| `GET /agents/{name}/versions/{v_from}/diff/{v_to}` | `agents.diff_versions` | Registry | `registry.diff_versions(...)` → `DiffResult` |
|
||||
| `GET /policies` | `policies.list_policies` | PolicyStore | `store.list_policies()` |
|
||||
| `GET /policies/{name}/versions` | `policies.list_versions` | PolicyStore | `store.list_versions(name)` (404) |
|
||||
| `GET /violations?trace_id=&severity=` | `violations.list_violations` | Settings | `read_violations(data_dir)` (lee `violations.jsonl`) + filtros opcionales por `trace_id` / `severity`. |
|
||||
| `GET /health` | (en `main.py`) | — | `{"status":"ok"}` |
|
||||
|
||||
Modelos de petición (en `executions.py`): `InvokeRequest{input: str, version: str | None}`, `ApproveRequest{approved_action_ids: list[str]=[], comment: str | None}`, `RejectRequest{reason: str}`. Modelos de respuesta = los de `domain/` (`response_model=…`).
|
||||
|
||||
Helpers internos de `executions.py` (índice de ejecuciones, en `data_dir/execution_index.json`):
|
||||
`_index_path`, `_load_index`, `_record_execution(data_dir, trace_id, agent_name, version)`,
|
||||
`_resolve(registry, policies, data_dir, trace_id) → (AgentDefinition, PolicyDefinition)` (con `HTTPException` 404/500),
|
||||
`_ensure_awaiting(orchestrator, agent_def, policy, trace_id)` (`HTTPException` 404/409).
|
||||
|
||||
---
|
||||
|
||||
## 5. Persistencia: quién escribe / lee qué
|
||||
|
||||
| Artefacto | Formato | Lo escribe | Lo lee |
|
||||
|-----------|---------|------------|--------|
|
||||
| `agents/<n>/versions/<v>.yaml`, `agents/<n>/index.yaml` | YAML | `FileSystemAgentRegistry.upsert_version` (a mano en este MVP) | `FileSystemAgentRegistry.get_*` / `list_*` / `diff_versions` |
|
||||
| `policies/<n>/versions/<v>.yaml`, `policies/<n>/index.yaml` | YAML | a mano | `FileSystemPolicyStore.get_policy` / `list_*` |
|
||||
| `data/checkpoints.sqlite` | SQLite (LangGraph `AsyncSqliteSaver`) | el grafo, en cada `ainvoke`/transición de nodo (incl. el `interrupt()`) — vía `runtime/checkpointer.build_checkpointer(data_dir)` | `orchestrator.snapshot` / `resume` (`graph.aget_state` / `Command(resume=…)`) |
|
||||
| `data/execution_index.json` | JSON `{trace_id: {agent_name, version}}` | `executions._record_execution` — en **cada** `invoke` | `executions._resolve`, `executions.list_executions` |
|
||||
| `data/executions.jsonl` | JSONL append-only (`AgentExecution` por línea) | `api/persistence.append_execution` — en `invoke` **si el status es terminal** y siempre tras `approve`/`reject` | `api/persistence.read_execution_summaries` (endpoint `GET /executions`) |
|
||||
| `data/violations.jsonl` | JSONL append-only (`GuardrailViolation` por línea) | `api/persistence.append_violation` — en `invoke`, una por violación de la ejecución | `api/persistence.read_violations` (endpoint `GET /violations`) |
|
||||
|
||||
Hashing/versionado: `registry/versioning.compute_hash(yaml_text)` = SHA-256 del YAML con espacios finales recortados; `unified_diff(a, b, label_a, label_b)` = `difflib.unified_diff`. (En los `index.yaml` de ejemplo el `hash` es `"pending"` y nadie lo valida al cargar.)
|
||||
|
||||
---
|
||||
|
||||
## 6. Dashboard ↔ Core (`agentforge_dashboard`)
|
||||
|
||||
El dashboard no comparte código con el core: solo lo llama por HTTP a través de
|
||||
`CoreClient` (`dashboard/src/agentforge_dashboard/client.py`, httpx síncrono con 2
|
||||
retries, `base_url = AGENTFORGE_CORE_URL`; mapea 404/409/422 → `{"error": <json>}`).
|
||||
|
||||
| `CoreClient.<método>` | Endpoint del core | Página(s) que lo usan |
|
||||
|-----------------------|-------------------|-----------------------|
|
||||
| `health()` | `GET /health` | `app.py` (sidebar) |
|
||||
| `list_agents()` | `GET /agents` | Registro, Ejecutar |
|
||||
| `get_agent(name)` | `GET /agents/{name}` | Registro |
|
||||
| `list_versions(name)` | `GET /agents/{name}/versions` | Registro |
|
||||
| `diff_versions(name, a, b)` | `GET /agents/{name}/versions/{a}/diff/{b}` | Registro (→ `components/diff_view`) |
|
||||
| `invoke_agent(name, body)` | `POST /agents/{name}/invoke` | Ejecutar |
|
||||
| `list_executions()` | `GET /executions` | Aprobaciones (filtra `status=="awaiting_approval"`), Historial |
|
||||
| `get_execution(trace_id)` | `GET /executions/{trace_id}` | Aprobaciones, Historial (→ `components/trace_view`, `violation_view`) |
|
||||
| `approve(trace_id, body)` | `POST /executions/{trace_id}/approve` | Aprobaciones |
|
||||
| `reject(trace_id, body)` | `POST /executions/{trace_id}/reject` | Aprobaciones |
|
||||
| `list_violations(**filters)` | `GET /violations?…` | Historial |
|
||||
| `list_policies()` | `GET /policies` | Politicas |
|
||||
| `list_policy_versions(name)` | `GET /policies/{name}/versions` | Politicas |
|
||||
|
||||
Páginas (Streamlit multipágina; el nº y el emoji del nombre del fichero son la
|
||||
navegación): `app.py` (raíz), `pages/1_🏛️_Registro.py`, `pages/2_▶️_Ejecutar.py`,
|
||||
`pages/3_🤝_Aprobaciones.py`, `pages/4_📜_Historial.py`, `pages/5_📐_Politicas.py`.
|
||||
Componentes reutilizables: `components/diff_view.render_unified_diff(diff_text)`,
|
||||
`components/trace_view.render_trace(decision_path)`, `components/violation_view.render_violations(violations)`.
|
||||
|
||||
---
|
||||
|
||||
## 7. Arranque y ciclo de vida
|
||||
|
||||
**Proceso core** (`uvicorn agentforge_core.main:app`):
|
||||
1. Import de `agentforge_core.main` ⇒ se ejecuta `app = create_app()`:
|
||||
`Settings()` → `configure_logging(level=settings.log_level)` → `FastAPI(...)` →
|
||||
`app.add_middleware(TraceIdMiddleware)` → registra `GET /health` →
|
||||
`from agentforge_core.api import agents, executions, policies, violations` →
|
||||
`include_router` ×5.
|
||||
2. Las dependencias (`deps.get_registry`, `get_policy_store`, `get_llm_provider`,
|
||||
`get_guardrail_engine`, `get_orchestrator`) **no** se construyen aún; se
|
||||
construyen y cachean en la **primera request** que las inyecta.
|
||||
3. Cada request: `TraceIdMiddleware.dispatch` → router → resuelve `Depends(...)` (que
|
||||
pueden disparar la construcción perezosa) → handler → respuesta con `X-Trace-Id`.
|
||||
|
||||
**Contenedores** (`docker-compose.yml`): servicio `core` (`core/Dockerfile`,
|
||||
`uvicorn agentforge_core.main:app --host 0.0.0.0 --port 8000`, `HEALTHCHECK` →
|
||||
`curl /health`, monta `./agents:ro`, `./policies:ro`, `./data:rw`, env
|
||||
`DATA_DIR=/app/data`, `AGENTS_DIR=/app/agents`, `POLICIES_DIR=/app/policies`,
|
||||
`env_file: .env`); servicio `dashboard` (`dashboard/Dockerfile`, `streamlit run
|
||||
app.py`, `HEALTHCHECK` → `/_stcore/health`, `depends_on: core: service_healthy`,
|
||||
env `AGENTFORGE_CORE_URL=http://core:8000`, monta `./agents:ro` para leer los
|
||||
`examples/*.txt`).
|
||||
|
||||
**`Settings` (env vars)** — `config.py`:
|
||||
|
||||
| Campo | Env var | Default | Lo consume |
|
||||
|-------|---------|---------|------------|
|
||||
| `llm_provider` | `LLM_PROVIDER` | `mock` | `build_llm_provider` |
|
||||
| `llm_fallback_provider` | `LLM_FALLBACK_PROVIDER` | `""` | (declarado; la factory aún no lo usa) |
|
||||
| `azure_openai_*` | `AZURE_OPENAI_*` | `""` / `2024-08-01-preview` | `AzureOpenAIProvider` |
|
||||
| `openai_api_key` / `openai_model` | `OPENAI_API_KEY` / `OPENAI_MODEL` | `""` / `gpt-4o` | `OpenAIProvider` |
|
||||
| `guardrails_nemo_enabled` | `GUARDRAILS_NEMO_ENABLED` | `False` | `build_guardrail_engine` |
|
||||
| `log_level` | `LOG_LEVEL` | `INFO` | `configure_logging` |
|
||||
| `data_dir` | `DATA_DIR` | `./data` | `build_checkpointer`, `_record_execution`, `append_*`, `read_*` |
|
||||
| `agents_dir` | `AGENTS_DIR` | `./agents` | `build_agent_registry` |
|
||||
| `policies_dir` | `POLICIES_DIR` | `./policies` | `build_policy_store` |
|
||||
|
||||
---
|
||||
|
||||
## 8. Aristas y "gotchas"
|
||||
|
||||
- **`llm_fallback_provider`**: existe en `Settings` y en `.env.example`, pero
|
||||
`build_llm_provider` aún no lo aplica (queda como punto de extensión).
|
||||
- **`build_llm_provider`** usa `match` sin `case _:`; al ser el tipo un `Literal`
|
||||
de tres valores es exhaustivo, pero un valor inesperado caería en "ninguna rama".
|
||||
- **NeMo**: `NeMoGuardrailsEngine` solo entra al `CompositeGuardrailEngine` si
|
||||
`GUARDRAILS_NEMO_ENABLED=true`; por defecto el composite tiene un solo engine.
|
||||
- **`detect_pii`** se comporta distinto según el entorno: con `presidio-analyzer`
|
||||
instalado (imagen Docker) usa Presidio + `en_core_web_sm`; sin él (venv local
|
||||
típico), regex fallback. La política `default` pide solo recognizers de patrón
|
||||
(`EMAIL_ADDRESS`, `ES_NIF`, `IP_ADDRESS`, `IBAN_CODE`) para evitar falsos
|
||||
positivos del NER.
|
||||
- **`AsyncSqliteSaver`** liga su conexión `aiosqlite` al event loop activo: por eso
|
||||
`build_checkpointer` es un *async context manager* y el orchestrator lo abre y
|
||||
cierra en cada operación (no se reusa entre llamadas). Requiere `aiosqlite<0.21`.
|
||||
- **`data/`** debe existir y ser escribible por el proceso/usuario del contenedor
|
||||
(`agent`, uid 1000); si no, `invoke` fallará al escribir el log/checkpoint.
|
||||
- **`make test` / `make lint` / `make smoke`** asumen que el venv está activado
|
||||
(los binarios `pytest`/`ruff`/`mypy`/`docker` en el `PATH`).
|
||||
- **`tests/integration/`** monta la app con `TestClient` apuntando a los assets
|
||||
*reales* del repo (`agents/`, `policies/`); `tests/unit/test_api_*` usan
|
||||
`tests/fixtures/` (versiones mínimas). Ambos hacen `deps.*.cache_clear()`.
|
||||
@@ -0,0 +1,542 @@
|
||||
# AgentForge explicado de principio a fin
|
||||
|
||||
> **Para quién es esto.** Una guía didáctica para alguien que llega nuevo al
|
||||
> proyecto —técnico o no— y quiere entender *qué hace*, *por qué está hecho así*
|
||||
> y *cómo encajan las piezas* sin tener que leer todo el código primero.
|
||||
>
|
||||
> Si solo quieres arrancarlo, ve al [`README.md`](../README.md). Si quieres las
|
||||
> decisiones técnicas en bruto, ve a [`ARCHITECTURE.md`](../ARCHITECTURE.md). Si
|
||||
> quieres la referencia de cableado a bajo nivel (módulos, firmas, grafos de
|
||||
> dependencias, cadenas de llamada), ve a [`docs/componentes.md`](componentes.md).
|
||||
> Este documento está en medio: cuenta la historia.
|
||||
|
||||
---
|
||||
|
||||
## 1. El problema, en una frase
|
||||
|
||||
> Poner agentes de IA en producción **sin una capa de gobierno** produce sistemas
|
||||
> opacos: prompts que cambian sin historial, validaciones inconsistentes, acciones
|
||||
> de alto impacto sin supervisión y ninguna auditoría de lo que decidió el agente.
|
||||
|
||||
**AgentForge es el "plano de control" que pones *delante* de tus agentes** antes de
|
||||
dejarlos tocar nada importante. No es un framework para *construir* agentes; es la
|
||||
capa que los **cataloga, versiona, valida, ejecuta de forma supervisada y audita**.
|
||||
|
||||
El caso de ejemplo que trae el repo es un agente de operaciones de telco
|
||||
(`incident_analyzer`): recibe la descripción de un incidente de plataforma de voz
|
||||
(caída de registros SIP, degradación de MOS, saturación de HSS...) y propone
|
||||
acciones con análisis de riesgo y plan de rollback. Acciones de riesgo alto quedan
|
||||
**pausadas esperando aprobación humana**. Todo queda registrado.
|
||||
|
||||
---
|
||||
|
||||
## 2. Las seis ideas grandes
|
||||
|
||||
Si entiendes estas seis ideas, entiendes el proyecto. Todo lo demás son detalles.
|
||||
|
||||
| # | Idea | Dónde vive |
|
||||
|---|------|------------|
|
||||
| 1 | **Agentes y políticas como ficheros declarativos, versionados como Git.** Un agente es un YAML (prompt, modelo, esquema de salida, umbral de aprobación). Cambias el YAML → nueva versión, con hash y diff. Sin redeploy. | `agents/`, `policies/`, `registry/` |
|
||||
| 2 | **Los guardrails son una *política*, no código disperso.** Una política lista validadores de entrada y de salida con su configuración. El motor los aplica; "qué se valida" es configuración. | `policies/`, `guardrails/` |
|
||||
| 3 | **La ejecución del agente es un grafo de estados con checkpoints.** No es "llama al LLM y ya"; es un flujo: validar entrada → razonar → validar salida → proponer acciones → puerta de aprobación → finalizar. Cada paso se persiste. | `runtime/` (LangGraph) |
|
||||
| 4 | **Human-in-the-Loop (HITL) de verdad.** Si el agente propone algo arriesgado, el grafo **se pausa** en mitad de la ejecución, el estado se guarda en disco, y se reanuda más tarde —incluso tras reiniciar el proceso— cuando un humano aprueba o rechaza. | nodo `approve_gate` + endpoints `/approve`, `/reject` |
|
||||
| 5 | **Trazabilidad obligatoria.** Cada petición lleva un `trace_id` (UUID) que se propaga por el middleware → los logs → la API → los ficheros de auditoría. Cada paso del grafo deja una entrada en el `decision_path` con su duración. | `observability/`, `domain/execution.py` |
|
||||
| 6 | **Todo lo "intercambiable" está detrás de una interfaz + un factory.** El proveedor de LLM, el motor de guardrails, el registry... son `Protocol`s con varias implementaciones. Un factory elige cuál según la configuración. Cambias `.env`, no el código. | `llm/`, `guardrails/`, `registry/` (los `factory.py`) |
|
||||
|
||||
---
|
||||
|
||||
## 3. Vista de pájaro: dos servicios
|
||||
|
||||
AgentForge son **dos procesos** que se hablan por HTTP/JSON:
|
||||
|
||||
```
|
||||
┌───────────────────────────── docker-compose ──────────────────────────────┐
|
||||
│ │
|
||||
│ ┌──────────────────────┐ HTTP/JSON ┌────────────────────────┐ │
|
||||
│ │ agentforge-dashboard │ ───────────────► │ agentforge-core │ │
|
||||
│ │ Streamlit :8501 │ ◄─────────────── │ FastAPI :8000 │ │
|
||||
│ │ (la "consola") │ │ (el cerebro) │ │
|
||||
│ └──────────────────────┘ └───────────┬────────────┘ │
|
||||
│ │ │
|
||||
│ ┌────────────────┬────────────────────┼───────────────┤
|
||||
│ ▼ ▼ ▼ │
|
||||
│ capa LLM capa Guardrails runtime LangGraph │
|
||||
│ (Strategy+factory) (Strategy+factory) (grafo + checkpointer) │
|
||||
│ │ │ │ │
|
||||
│ └────────────────┴──────────┬──────────┘ │
|
||||
│ ▼ │
|
||||
│ Persistencia: YAML · JSON · JSONL · SQLite │
|
||||
└────────────────────────────────────────────────────────────────────────────┘
|
||||
```
|
||||
|
||||
- **`agentforge-core`** (FastAPI, puerto 8000) — todo el dominio: registry de
|
||||
agentes, motor de guardrails, runtime de ejecución, persistencia. No tiene UI.
|
||||
- **`agentforge-dashboard`** (Streamlit, puerto 8501) — una consola visual. **No
|
||||
contiene lógica de negocio**: es un cliente HTTP del core con cinco páginas.
|
||||
|
||||
La separación importa: el core podría servir a una CLI, a otro servicio, a un
|
||||
pipeline... el dashboard es solo una de las caras posibles.
|
||||
|
||||
### Las capas del core (de fuera hacia dentro)
|
||||
|
||||
```
|
||||
HTTP ─► middlewares ─► routers (api/) ─► dependencias (deps.py)
|
||||
│ │ ← aquí se inyectan los
|
||||
│ │ objetos del dominio
|
||||
▼ ▼
|
||||
orchestrator ──► graph (LangGraph) ──► nodes
|
||||
│ │
|
||||
│ ├─► llm/ (¿qué dice el LLM?)
|
||||
│ ├─► guardrails/ (¿pasa los filtros?)
|
||||
│ └─► domain/ (¿qué forma tienen los datos?)
|
||||
▼
|
||||
checkpointer (SQLite) + persistence (JSONL) + registry (YAML/JSON)
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 4. Recorrido por los módulos (y quién depende de quién)
|
||||
|
||||
El código del core vive bajo `core/src/agentforge_core/`. Lo agrupo por capas, de
|
||||
las más internas (sin dependencias) a las más externas.
|
||||
|
||||
### 4.1 `domain/` — el vocabulario del sistema
|
||||
|
||||
Modelos Pydantic puros. **No dependen de nada del proyecto**; todo lo demás depende
|
||||
de ellos. Son el "idioma común".
|
||||
|
||||
| Fichero | Qué define |
|
||||
|---------|------------|
|
||||
| `agent.py` | `AgentDefinition` (prompt, modelo, `output_schema`, `guardrails`, `risk_threshold_for_hitl`, ...), `LLMConfig`, `AgentVersionMeta`. |
|
||||
| `policy.py` | `PolicyDefinition` (listas de `PolicyValidator` de entrada y de salida, `on_validator_error`), `PolicyVersionMeta`. |
|
||||
| `guardrail.py` | `GuardrailViolation` (trace_id, stage `input`/`output`, validator, severity, message, blocked). |
|
||||
| `execution.py` | `AgentExecution` (el "expediente" de una ejecución: status, `decision_path`, `violations`, `proposed_actions`, `needs_human_for`, `final_output`, `error`), `ProposedAction`, `DecisionStep`, `AgentExecutionSummary`. |
|
||||
|
||||
> **Pista didáctica:** si quieres entender el sistema rápido, empieza leyendo
|
||||
> `domain/execution.py`. Te dice exactamente qué información se produce y se guarda.
|
||||
|
||||
### 4.2 `config.py` y `observability/logging.py` — los cimientos transversales
|
||||
|
||||
- **`config.py` → `Settings`** (pydantic-settings): lee `.env` (proveedor LLM,
|
||||
rutas de `agents/`/`policies/`/`data/`, nivel de log, claves de Azure/OpenAI,
|
||||
flags de NeMo). Es el único sitio que sabe de variables de entorno. **Todos los
|
||||
`factory.py` reciben un `Settings`.**
|
||||
- **`observability/logging.py`**: configura `structlog` con salida JSON y un
|
||||
contexto donde se "bind-ea" el `trace_id`. Lo usa todo el código que loguea
|
||||
(`log = structlog.get_logger(__name__)`).
|
||||
|
||||
Dependen: de nada del proyecto. Dependen de ellos: prácticamente todo.
|
||||
|
||||
### 4.3 `llm/` — la capa de proveedores de LLM (Strategy pattern)
|
||||
|
||||
| Fichero | Rol |
|
||||
|---------|-----|
|
||||
| `base.py` | El `Protocol` `LLMProvider` (`async complete(messages, temperature, max_tokens) -> CompletionResult`) y los tipos `Message`, `CompletionResult`. **La interfaz.** |
|
||||
| `mock.py` | `MockProvider`: determinista, sin claves. Elige una respuesta canónica buscando subcadenas (`"sip"`, `"mos"`, `"hss"`) en el input. Es lo que hace que el demo funcione out-of-the-box. |
|
||||
| `azure.py`, `openai.py` | Implementaciones reales (Azure OpenAI / OpenAI) con reintentos exponenciales y, opcionalmente, *fallback* entre proveedores. |
|
||||
| `factory.py` | `build_llm_provider(settings)`: devuelve el provider según `LLM_PROVIDER` (`mock` / `azure` / `openai`), y envuelve un *fallback* opcional. |
|
||||
|
||||
Depende de: `domain/`, `config.py`. Dependen de él: el `runtime/` (el nodo
|
||||
`llm_reason`) y el `deps.py` de la API.
|
||||
|
||||
### 4.4 `registry/` — catálogo y versionado de agentes y políticas
|
||||
|
||||
| Fichero | Rol |
|
||||
|---------|-----|
|
||||
| `repository.py` | `FileSystemAgentRegistry`: lee `agents/<nombre>/index.yaml` + `versions/<id>.yaml`. CRUD de agentes y versiones, lectura de la versión "activa". |
|
||||
| `policy_store.py` | `FileSystemPolicyStore`: lo mismo para `policies/`. Devuelve `PolicyDefinition`s. |
|
||||
| `versioning.py` | `compute_hash(yaml_text)` (SHA-256 sobre el contenido normalizado) y el *diff unificado* entre dos versiones. Es la "magia tipo Git": versiones inmutables identificadas por hash, comparables. |
|
||||
| `factory.py` | `build_registry(settings)` y `build_policy_store(settings)`. |
|
||||
|
||||
Depende de: `domain/`, `config.py`. Dependen de él: la API (`/agents`, `/policies`),
|
||||
el `orchestrator` (para resolver qué agente/política ejecutó una traza dada).
|
||||
|
||||
### 4.5 `guardrails/` — el motor de validación
|
||||
|
||||
Aquí está el corazón del "gobierno". La estructura sigue otra vez **Protocol +
|
||||
implementaciones + composite + factory**:
|
||||
|
||||
```
|
||||
GuardrailEngine (Protocol) ← base.py: validate_input / validate_output
|
||||
▲
|
||||
┌────────────┼─────────────┐
|
||||
│ │
|
||||
GuardrailsAIEngine NeMoGuardrailsEngine ← guardrails_ai.py / nemo.py
|
||||
(el real: aplica los (stub; desactivado
|
||||
validadores de la por defecto)
|
||||
política)
|
||||
│
|
||||
└──────────┐
|
||||
▼
|
||||
CompositeGuardrailEngine ← composite.py: corre N sub-engines
|
||||
(corre los sub-engines en paralelo en paralelo (asyncio.gather) y
|
||||
con asyncio y agrega las violaciones) une las listas de violaciones
|
||||
|
||||
factory.py: build_guardrail_engine(settings) → ensambla el Composite
|
||||
validators.py: las funciones concretas — detect_pii, prompt_injection,
|
||||
toxic_language, forbidden_topics, schema_match, pii_leakage,
|
||||
forbidden_action_keywords, telco_safety_rules
|
||||
```
|
||||
|
||||
**Cómo se conecta una política con un validador:** la política
|
||||
(`policies/default/versions/v1.yaml`) lista, por ejemplo:
|
||||
|
||||
```yaml
|
||||
input_validators:
|
||||
- type: detect_pii
|
||||
config: { entities: [EMAIL_ADDRESS, ES_NIF, IP_ADDRESS, IBAN_CODE], severity_on_match: block }
|
||||
- type: prompt_injection
|
||||
config: { severity_on_match: block }
|
||||
...
|
||||
on_validator_error: fail_closed
|
||||
```
|
||||
|
||||
El `GuardrailsAIEngine` recorre esa lista, busca cada `type` en su registro de
|
||||
funciones de `validators.py`, la llama con el texto y el `config`, y junta las
|
||||
`GuardrailViolation` que devuelva. `on_validator_error: fail_closed` significa que
|
||||
**si un validador peta, cuenta como bloqueo** (seguridad antes que disponibilidad).
|
||||
|
||||
> **Sobre `detect_pii`:** usa Presidio (con el modelo spaCy `en_core_web_sm`) si
|
||||
> está instalado; si no, cae a una detección por regex (email, teléfono, NIF
|
||||
> español, IP). Por eso el comportamiento en local (sin Presidio) y en Docker
|
||||
> (con Presidio) puede diferir — la política solo pide *recognizers de patrón*
|
||||
> fiables para evitar falsos positivos del NER.
|
||||
|
||||
Depende de: `domain/`, `config.py`. Dependen de él: el `runtime/` (los nodos
|
||||
`validate_input`/`validate_output`) y el `deps.py`.
|
||||
|
||||
### 4.6 `runtime/` — la ejecución como grafo de estados (LangGraph)
|
||||
|
||||
Esta es la capa más "viva". Modela una ejecución del agente como un grafo dirigido.
|
||||
|
||||
| Fichero | Rol |
|
||||
|---------|-----|
|
||||
| `state.py` | `AgentState`: un `TypedDict` con todo lo que fluye por el grafo (input, salida del LLM, acciones propuestas, violaciones acumuladas, `decision_path`, status, decisión humana, ...). Algunos campos usan reducers (`Annotated[list, operator.add]`) para que cada nodo *añada* en vez de sobrescribir. |
|
||||
| `nodes.py` | Las **funciones-nodo**, construidas por *factories* parametrizadas con el agente, la política, el motor de guardrails y el provider LLM: `validate_input`, `llm_reason`, `validate_output`, `propose_actions`, `approve_gate`, `finalize`. Cada nodo añade un `DecisionStep` con su duración. |
|
||||
| `graph.py` | `build_graph(...)`: cablea los nodos y las **aristas condicionales** (p. ej. si `validate_input` bloqueó → salta directo al final). Compila el grafo con un `checkpointer`. |
|
||||
| `checkpointer.py` | `build_checkpointer(data_dir)`: un `AsyncSqliteSaver` de LangGraph sobre `data_dir/checkpoints.sqlite`, expuesto como *context manager* asíncrono. Es lo que hace que un `awaiting_approval` **sobreviva a un reinicio**. |
|
||||
| `orchestrator.py` | `AgentOrchestrator`: la **única puerta de entrada** al runtime. Tres operaciones: `invoke()` (lanza), `resume()` (reanuda un HITL con la decisión), `snapshot()` (lee el estado actual sin avanzarlo). Cada llamada abre su propio checkpointer y traduce el `StateSnapshot` de LangGraph a un `AgentExecution` del dominio. |
|
||||
|
||||
**El grafo, dibujado:**
|
||||
|
||||
```
|
||||
START
|
||||
│
|
||||
▼
|
||||
┌─────────────────┐ bloqueada (PII, injection...)
|
||||
│ validate_input │ ─────────────────────────────────────────► END
|
||||
└────────┬────────┘
|
||||
│ ok
|
||||
▼
|
||||
┌─────────────────┐ LLM no disponible / error
|
||||
│ llm_reason │ ─────────────────────────────────────────► END
|
||||
└────────┬────────┘
|
||||
│ ok
|
||||
▼
|
||||
┌─────────────────┐ salida no cumple el esquema / PII en la salida
|
||||
│ validate_output │ ─────────────────────────────────────────► END
|
||||
└────────┬────────┘
|
||||
│ ok
|
||||
▼
|
||||
┌─────────────────┐
|
||||
│ propose_actions │ (extrae las acciones propuestas del JSON del LLM)
|
||||
└────────┬────────┘
|
||||
│
|
||||
▼
|
||||
┌─────────────────┐ ¿hay alguna acción con risk_score ≥ umbral
|
||||
│ approve_gate │ o requires_approval=True?
|
||||
└────────┬────────┘
|
||||
│
|
||||
┌────┴───────────────────────────────┐
|
||||
│ no │ sí
|
||||
▼ ▼
|
||||
┌──────────┐ interrupt({...}) ──► el grafo SE PAUSA aquí.
|
||||
│ finalize │ El estado queda en checkpoints.sqlite.
|
||||
└────┬─────┘ Más tarde llega resume(decision={...})
|
||||
│ y se reanuda en este mismo punto.
|
||||
▼ │
|
||||
END ▼
|
||||
┌──────────┐
|
||||
│ finalize │ (filtra a las acciones aprobadas;
|
||||
└────┬─────┘ si fue rechazo → status=failed)
|
||||
▼
|
||||
END
|
||||
```
|
||||
|
||||
Depende de: `llm/`, `guardrails/`, `domain/`, `config.py` (vía los objetos que le
|
||||
inyectan). Dependen de él: la API (`api/executions.py` solo conoce el
|
||||
`AgentOrchestrator`, no LangGraph).
|
||||
|
||||
### 4.7 `api/` — la fachada HTTP (FastAPI)
|
||||
|
||||
| Fichero | Rol |
|
||||
|---------|-----|
|
||||
| `middlewares.py` | `TraceIdMiddleware`: lee `X-Trace-Id` de la petición (o genera uno), lo bind-ea al contexto de structlog, y lo devuelve en la respuesta. **Es el origen del hilo de trazabilidad.** |
|
||||
| `deps.py` | Las **dependencias inyectables**. Cada `get_*` (`get_settings`, `get_registry`, `get_policy_store`, `get_llm_provider`, `get_guardrail_engine`, `get_orchestrator`) está cacheada con `lru_cache(maxsize=1)`: la app construye cada cosa **una sola vez**. Expone aliases (`RegistryDep = Annotated[..., Depends(get_registry)]`, etc.) que los routers piden por parámetro. Los tests llaman a `.cache_clear()` para reconstruir todo apuntando a un `DATA_DIR` temporal. |
|
||||
| `persistence.py` | Helpers de **log append-only en JSONL**: `append_execution`, `append_violation`, `read_execution_summaries`. Inmutable, auditable, fácil de "shipear" a un sistema de logs. |
|
||||
| `agents.py` | Router `/agents`: list, get, versiones, `GET /agents/{n}/versions/{a}/diff/{b}`. |
|
||||
| `executions.py` | El más cargado: `POST /agents/{n}/invoke`, `GET /executions`, `GET /executions/{trace_id}`, `POST /executions/{trace_id}/approve`, `POST /executions/{trace_id}/reject`. Mantiene además `execution_index.json` (mapa `trace_id → agente/versión`) para poder reanudar tras un reinicio. |
|
||||
| `policies.py`, `violations.py` | Routers `/policies` y `/violations` (este con filtros por severidad, stage, etc.). |
|
||||
|
||||
Y la **raíz de la app**: `core/src/agentforge_core/main.py` → `create_app()` instancia
|
||||
`FastAPI`, añade el `TraceIdMiddleware`, monta los routers, y expone `/health`.
|
||||
|
||||
Depende de: todo lo de arriba (vía `deps.py`). Dependen de él: el dashboard (por
|
||||
HTTP) y los tests de integración (`tests/integration/`, vía `TestClient`).
|
||||
|
||||
### 4.8 `dashboard/` — la consola Streamlit
|
||||
|
||||
| Fichero | Rol |
|
||||
|---------|-----|
|
||||
| `client.py` | `CoreClient`: un cliente HTTP síncrono (httpx, con reintentos) que envuelve **todos** los endpoints del core y mapea 404/409/422 a `{"error": ...}`. Es lo único que sabe hablar con el core. |
|
||||
| `app.py` | La página raíz: sidebar de branding + un health-check del core. |
|
||||
| `pages/1_🏛️_Registro.py` | Catálogo de agentes: detalle (prompt, esquema, LLM, guardrails), tabla de versiones, **diff coloreado v1↔v2**. |
|
||||
| `pages/2_▶️_Ejecutar.py` | Lanza un agente: botones con los escenarios pregrabados, textarea, `invoke`, y render del status + output + violaciones + *timeline* del `decision_path`. Avisa si quedó en `awaiting_approval`. |
|
||||
| `pages/3_🤝_Aprobaciones.py` | La cola de HITL: lista las ejecuciones `awaiting_approval`, muestra cada acción propuesta (risk_score coloreado, target, rollback_plan) con un checkbox, y aprueba el subconjunto elegido o rechaza con motivo. |
|
||||
| `pages/4_📜_Historial.py` | Pestaña de ejecuciones (tabla + detalle por `trace_id`) y pestaña de violaciones (filtrable por severidad). |
|
||||
| `pages/5_📐_Politicas.py` | Inventario de políticas: validadores de entrada/salida (cada uno expandible con su config) y versiones. |
|
||||
| `components/` | Trozos reutilizables de UI: `diff_view` (pinta `+`/`-`/`@@`), `trace_view` (el timeline del `decision_path`), `violation_view` (badges de severidad). |
|
||||
|
||||
Depende de: el `agentforge-core` por HTTP (vía `AGENTFORGE_CORE_URL`). **Nadie del
|
||||
core depende del dashboard.** No tiene tests unitarios (mal coste/beneficio para
|
||||
Streamlit); su verificación es la checklist manual de [`docs/manual_qa.md`](manual_qa.md).
|
||||
|
||||
### 4.9 Lo que no es código: `agents/`, `policies/`, `data/`
|
||||
|
||||
- **`agents/incident_analyzer/`** — el agente de ejemplo: `index.yaml` (catálogo de
|
||||
versiones), `versions/v1.yaml` y `v2.yaml`, y `examples/*.txt` (tres escenarios de
|
||||
incidente de telco que el demo usa).
|
||||
- **`policies/default/`** — la política de guardrails de ejemplo (`index.yaml` +
|
||||
`versions/v1.yaml`).
|
||||
- **`data/`** — estado *runtime* (gitignored): `checkpoints.sqlite`,
|
||||
`executions.jsonl`, `violations.jsonl`, `execution_index.json`. Se crea sola.
|
||||
|
||||
---
|
||||
|
||||
## 5. El patrón que se repite: `Protocol` + `factory` + `Settings`
|
||||
|
||||
Tres veces (LLM, guardrails, registry) verás la misma estructura:
|
||||
|
||||
```
|
||||
base.py → un Protocol (la interfaz: "qué se puede hacer")
|
||||
<impl_a>.py → una implementación (p. ej. mock)
|
||||
<impl_b>.py → otra implementación (p. ej. azure)
|
||||
factory.py → build_X(settings: Settings) -> X ← elige y monta
|
||||
```
|
||||
|
||||
¿Por qué? Porque permite **cambiar el comportamiento sin tocar el código**: pones
|
||||
`LLM_PROVIDER=azure` en `.env` y el factory te da el provider de Azure; pones
|
||||
`mock` y tienes un demo determinista sin claves. Lo mismo con guardrails (puedes
|
||||
añadir el engine de NeMo) y con el registry (hoy es de ficheros; mañana podría ser
|
||||
de base de datos). Es el principio de **"configuración antes que código"**.
|
||||
|
||||
Y todo se enchufa **una sola vez** al arrancar, en `api/deps.py` (los `lru_cache`):
|
||||
el registry, el policy store, el provider, el motor de guardrails y el orchestrator
|
||||
son singletons del proceso. Los routers solo los *piden*; no saben construirlos.
|
||||
|
||||
---
|
||||
|
||||
## 6. El recorrido de UNA petición, de principio a fin
|
||||
|
||||
Esta es la sección que conviene leer despacio: aquí se ve cómo encaja todo. Sigamos
|
||||
`POST /agents/incident_analyzer/invoke` con el escenario SIP (el que dispara HITL).
|
||||
|
||||
### Acto 1 — la petición entra y se ejecuta el grafo
|
||||
|
||||
```
|
||||
Cliente (dashboard o curl)
|
||||
│ POST /agents/incident_analyzer/invoke { "input": "<texto del incidente SIP>" }
|
||||
▼
|
||||
TraceIdMiddleware ............ genera trace_id = UUID, lo bind-ea al log
|
||||
▼
|
||||
router invoke_agent (api/executions.py)
|
||||
│ pide por inyección: RegistryDep, PolicyStoreDep, OrchestratorDep, SettingsDep
|
||||
│ registry.get_agent("incident_analyzer") → AgentDefinition (versión activa = v2)
|
||||
│ policies.get_policy(agent.guardrails[0]) → PolicyDefinition ("default")
|
||||
▼
|
||||
orchestrator.invoke(agent_def, policy, user_input)
|
||||
│ abre el checkpointer (AsyncSqliteSaver sobre data/checkpoints.sqlite)
|
||||
│ build_graph(agent_def, policy, provider, engine, checkpointer)
|
||||
│ graph.ainvoke(estado_inicial, config={thread_id: trace_id})
|
||||
▼
|
||||
┌──── el grafo ────────────────────────────────────────────────────────────┐
|
||||
│ validate_input → engine.validate_input(texto, policy, trace_id) │
|
||||
│ recorre los input_validators de la política en paralelo │
|
||||
│ (detect_pii, prompt_injection, ...) → 0 violaciones │
|
||||
│ añade DecisionStep("validate_input", ...) │
|
||||
│ llm_reason → provider.complete([system_prompt, user_input], ...) │
|
||||
│ (MockProvider ve "sip" → respuesta canónica SIP) │
|
||||
│ añade DecisionStep("llm_reason", model, tokens, ...) │
|
||||
│ validate_output → parsea el JSON del LLM; engine.validate_output(...) │
|
||||
│ (schema_match, pii_leakage, forbidden_action_keywords, │
|
||||
│ telco_safety_rules) → 0 violaciones │
|
||||
│ propose_actions → extrae proposed_actions del JSON: [act-1 (risk 4, │
|
||||
│ requires_approval), act-2 (risk 3)] │
|
||||
│ approve_gate → ¿alguna acción con risk ≥ 4 (umbral del agente) o │
|
||||
│ requires_approval? SÍ (act-1) → interrupt({...}) │
|
||||
│ ⇒ el grafo SE DETIENE. El estado se escribe en SQLite. │
|
||||
└──────────────────────────────────────────────────────────────────────────┘
|
||||
▼
|
||||
orchestrator._snapshot(...) → LangGraph reporta "hay un nodo pendiente"
|
||||
⇒ status = "awaiting_approval"
|
||||
⇒ needs_human_for = [act-1] (y devuelve un AgentExecution)
|
||||
▼
|
||||
router: status no es terminal → NO se escribe en executions.jsonl,
|
||||
pero SÍ se registra en execution_index.json: { trace_id → (incident_analyzer, v2) }
|
||||
▼
|
||||
respuesta 200 { status: "awaiting_approval", trace_id, needs_human_for: [act-1], decision_path: [...], ... }
|
||||
```
|
||||
|
||||
En el dashboard, la página **Ejecutar** muestra el timeline y un aviso "ve a
|
||||
Aprobaciones". La página **Aprobaciones** hace `GET /executions`, encuentra esta
|
||||
ejecución (el endpoint reconstruye las `awaiting_approval` desde el índice + el
|
||||
checkpointer) y muestra `act-1` con su risk_score, target y rollback_plan.
|
||||
|
||||
### Acto 2 — el humano decide; el grafo se reanuda
|
||||
|
||||
```
|
||||
Humano (en la página Aprobaciones, o curl)
|
||||
│ POST /executions/<trace_id>/approve { approved_action_ids: ["act-1"], comment: "ok rollback" }
|
||||
▼
|
||||
router approve_execution (api/executions.py)
|
||||
│ _resolve(...) → lee execution_index.json → sabe que fue (incident_analyzer, v2)
|
||||
│ reconstruye el AgentDefinition y la PolicyDefinition
|
||||
│ _ensure_awaiting(...) → orchestrator.snapshot(...) confirma que sigue en awaiting_approval
|
||||
▼
|
||||
orchestrator.resume(agent_def, policy, trace_id, decision={approved_action_ids:["act-1"], rejected:false})
|
||||
│ abre OTRA VEZ el checkpointer (mismo data_dir) — el estado pausado sigue ahí,
|
||||
│ aunque hubiera habido un reinicio del proceso entremedias
|
||||
│ graph.ainvoke(Command(resume=decision), config={thread_id: trace_id})
|
||||
▼
|
||||
┌──── el grafo continúa desde donde se quedó ──────────────────────────────┐
|
||||
│ approve_gate → el interrupt() devuelve la `decision` del humano │
|
||||
│ finalize → final_actions = acciones cuyo id está aprobado = [act-1] │
|
||||
│ final_output = { ...salida del LLM..., approved_actions:[act-1] }│
|
||||
│ status = "completed" │
|
||||
└──────────────────────────────────────────────────────────────────────────┘
|
||||
▼
|
||||
router: status terminal → append_execution(...) escribe el AgentExecution en executions.jsonl
|
||||
▼
|
||||
respuesta 200 { status: "completed", final_output: { ..., approved_actions: [act-1] }, ... }
|
||||
```
|
||||
|
||||
(Si en vez de `approve` se llama a `reject`, el nodo `finalize` ve `rejected: true` →
|
||||
`status = "failed"`, `error = "rejected_by_human"`. También terminal → al JSONL.)
|
||||
|
||||
### El mismo recorrido, como diagrama de secuencia
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant U as Cliente / Dashboard
|
||||
participant MW as TraceIdMiddleware
|
||||
participant API as router (api/executions.py)
|
||||
participant ORCH as AgentOrchestrator
|
||||
participant G as Grafo LangGraph
|
||||
participant CP as checkpoints.sqlite
|
||||
participant LOG as executions.jsonl / index.json
|
||||
|
||||
U->>MW: POST /agents/incident_analyzer/invoke {input}
|
||||
MW->>API: + trace_id
|
||||
API->>API: registry.get_agent · policies.get_policy
|
||||
API->>ORCH: invoke(agent, policy, input)
|
||||
ORCH->>G: ainvoke(estado, thread_id=trace_id)
|
||||
G->>G: validate_input → llm_reason → validate_output → propose_actions
|
||||
G->>G: approve_gate: hay riesgo alto → interrupt()
|
||||
G->>CP: persiste estado pausado
|
||||
ORCH-->>API: AgentExecution(status=awaiting_approval, needs_human_for=[act-1])
|
||||
API->>LOG: index.json[trace_id] = (incident_analyzer, v2)
|
||||
API-->>U: 200 {status: awaiting_approval, trace_id, ...}
|
||||
|
||||
Note over U,CP: ...más tarde (incluso tras reiniciar el core)...
|
||||
|
||||
U->>API: POST /executions/{trace_id}/approve {approved_action_ids:[act-1]}
|
||||
API->>LOG: lee index.json → (incident_analyzer, v2)
|
||||
API->>ORCH: resume(agent, policy, trace_id, decision)
|
||||
ORCH->>CP: reabre checkpointer (estado pausado sigue ahí)
|
||||
ORCH->>G: ainvoke(Command(resume=decision), thread_id=trace_id)
|
||||
G->>G: approve_gate (recibe decisión) → finalize → status=completed
|
||||
ORCH-->>API: AgentExecution(status=completed, final_output)
|
||||
API->>LOG: append a executions.jsonl
|
||||
API-->>U: 200 {status: completed, final_output, ...}
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 7. Persistencia: cuatro formas, cuatro razones
|
||||
|
||||
AgentForge no usa una sola base de datos; usa la herramienta adecuada para cada cosa.
|
||||
|
||||
| Qué | Cómo | Por qué así |
|
||||
|-----|------|-------------|
|
||||
| Definiciones de agentes y políticas | **YAML** por versión + `index.yaml` | Lo escribe un humano (legible, comentable) y lo cataloga la máquina. |
|
||||
| Identidad de versiones / comparación | **hash SHA-256** del contenido + `difflib` | Versiones inmutables identificables y comparables — "tipo Git". |
|
||||
| Mapa `trace_id → agente/versión` | **JSON** (`execution_index.json`) | Necesario para reanudar un HITL sabiendo qué configuración lo ejecutó; debe sobrevivir a reinicios. |
|
||||
| Log de ejecuciones terminales y de violaciones | **JSONL append-only** | Inmutable, auditable, trivial de "shipear" a un sistema de logs. |
|
||||
| Estado intermedio del grafo (incl. pausas HITL) | **SQLite** (`checkpoints.sqlite`, vía LangGraph) | Es lo que LangGraph espera; permite reanudar una ejecución pausada entre reinicios del proceso. |
|
||||
|
||||
---
|
||||
|
||||
## 8. Glosario
|
||||
|
||||
| Término | Significado en este proyecto |
|
||||
|---------|------------------------------|
|
||||
| **Agente** | Una configuración declarativa (YAML): prompt de sistema, modelo LLM, esquema de salida, lista de políticas de guardrails, umbral de riesgo para HITL. No es código. |
|
||||
| **Política (de guardrails)** | Un YAML con la lista de validadores de entrada y de salida y su configuración, más `on_validator_error` (`fail_closed`/`fail_open`). |
|
||||
| **Validador / guardrail** | Una función que mira el texto de entrada o la salida del LLM y devuelve cero o más `GuardrailViolation`. Ej.: `detect_pii`, `prompt_injection`, `schema_match`, `telco_safety_rules`. |
|
||||
| **Violación** | El resultado de un validador: `stage` (input/output), `validator`, `severity` (`info`/`warning`/`block`), `message`, `blocked`. |
|
||||
| **Ejecución (`AgentExecution`)** | El "expediente" de una invocación: `trace_id`, status, `decision_path`, `violations`, `proposed_actions`, `needs_human_for`, `final_output`, `error`. |
|
||||
| **`decision_path`** | La lista de pasos por los que pasó el grafo, cada uno con su `step`, `timestamp`, `duration_ms` y un `detail`. La "caja negra" de la ejecución. |
|
||||
| **HITL (Human-in-the-Loop)** | El patrón de pausar la ejecución cuando hay una acción arriesgada y esperar a que un humano apruebe o rechace. Implementado con `interrupt()` de LangGraph. |
|
||||
| **Checkpointer** | El componente de LangGraph que persiste el estado del grafo (aquí, en SQLite). Lo que hace posible reanudar un HITL tras un reinicio. |
|
||||
| **`trace_id`** | UUID que identifica una petición/ejecución y se propaga por middleware → logs → API → ficheros de auditoría. |
|
||||
| **Orchestrator** | La única clase que habla con LangGraph (`invoke`/`resume`/`snapshot`); aísla al resto del sistema del runtime. |
|
||||
| **Factory** | Función `build_X(settings)` que monta la implementación de `X` adecuada según la configuración. |
|
||||
| **Status de ejecución** | `running`, `awaiting_approval` (pausada en HITL), `blocked_by_guardrail` (la frenó un validador), `completed`, `failed`. |
|
||||
| **Estados de un agente** | `draft`, `active`, `deprecated` (metadato en su YAML). |
|
||||
|
||||
---
|
||||
|
||||
## 9. Mapa del repositorio y por dónde empezar a leer
|
||||
|
||||
```
|
||||
agentforge/
|
||||
├── core/
|
||||
│ ├── src/agentforge_core/
|
||||
│ │ ├── domain/ ← los modelos de datos (empieza por execution.py)
|
||||
│ │ ├── config.py · observability/ ← cimientos transversales
|
||||
│ │ ├── llm/ ← proveedores LLM (Protocol + impls + factory)
|
||||
│ │ ├── registry/ ← catálogo y versionado (YAML/JSON, hash, diff)
|
||||
│ │ ├── guardrails/ ← motor de validación (Protocol + composite + validators)
|
||||
│ │ ├── runtime/ ← el grafo LangGraph (state, nodes, graph, checkpointer, orchestrator)
|
||||
│ │ ├── api/ ← FastAPI (middlewares, deps, routers, persistence)
|
||||
│ │ └── main.py ← create_app(): ensambla la app
|
||||
│ ├── Dockerfile · requirements.txt
|
||||
├── dashboard/
|
||||
│ ├── src/agentforge_dashboard/ ← Streamlit (client + app + pages + components)
|
||||
│ ├── Dockerfile · requirements.txt
|
||||
├── agents/incident_analyzer/ ← el agente de ejemplo (YAMLs + escenarios .txt)
|
||||
├── policies/default/ ← la política de guardrails de ejemplo
|
||||
├── data/ ← estado runtime (gitignored)
|
||||
├── tests/ ← pytest: tests/unit/ y tests/integration/ (88 tests en total; el demo Streamlit no se testea con unit tests)
|
||||
├── docs/ ← este documento, manual_qa.md, futuro.md, superpowers/
|
||||
├── docker-compose.yml ← levanta core + dashboard
|
||||
├── Makefile ← install / test / test-all / lint / smoke / up / down
|
||||
├── README.md · ARCHITECTURE.md
|
||||
```
|
||||
|
||||
**Ruta de lectura sugerida (1 hora):**
|
||||
1. `domain/execution.py` y `domain/agent.py` — qué datos hay.
|
||||
2. `policies/default/versions/v1.yaml` y `agents/incident_analyzer/versions/v2.yaml` — cómo se declara todo.
|
||||
3. `runtime/nodes.py` y `runtime/graph.py` — el flujo de ejecución.
|
||||
4. `api/executions.py` — cómo se expone (invoke / approve / reject).
|
||||
5. `tests/integration/test_invoke_hitl.py` — el ciclo completo, en ~40 líneas.
|
||||
6. Levanta `docker compose up` y haz clic por las cinco páginas del dashboard.
|
||||
|
||||
---
|
||||
|
||||
## 10. Lo que aún no hace (a propósito)
|
||||
|
||||
Es un MVP. La detección de PII más fina (modelos grandes, español), autenticación,
|
||||
OpenTelemetry, multi-tenant, persistencia en Postgres, evaluadores LLM-as-judge,
|
||||
una cola de aprobaciones con SLA... están en el roadmap: [`docs/futuro.md`](futuro.md).
|
||||
El estado actual y cómo verificarlo: [`README.md`](../README.md) y [`docs/manual_qa.md`](manual_qa.md).
|
||||
@@ -5572,6 +5572,7 @@ git commit -m "feat(dashboard): página Politicas con detalle de validadores y v
|
||||
> **Nota de implementación (desviación del plan):**
|
||||
> - **Task 36:** ambos Dockerfiles pasan `docker build --check` sin warnings (se usó como lint previo).
|
||||
> - **Task 37 (smoke):** `docker compose build` falló la primera vez en `RUN python -m spacy download en_core_web_sm` del core — `guardrails-ai<0.6` arrastra `typer 0.12.x`, incompatible con el `click 8.3.x` que resolvió pip (`TypeError: Secondary flag is not valid for non-boolean flag`). Se mantuvo el `spacy download` del plan y se fijó `click>=8.1,<8.2` en `core/requirements.txt` (commit `fix(deps): fija click<8.2 …`); el rebuild funcionó. Smoke OK: `docker compose up -d` → `agentforge-core` *healthy* con `GET /health → 200 {"status":"ok"}`, `agentforge-dashboard` con `/_stcore/health → ok`, sin errores en logs; `docker compose down` limpio.
|
||||
> - **Walkthrough manual del dashboard (Task 38, ejecutado contra el stack):** al recorrer los flujos de `docs/manual_qa.md` vía la API del core aparecieron tres bugs (que los tests no veían porque el venv local no tiene Presidio y usa la rama de regex): (1) `detect_pii` instanciaba `AnalyzerEngine()` sin config → Presidio cargaba su modelo por defecto `en_core_web_lg`, que no está en la imagen → fail-closed bloqueaba todo → fix `fix(guardrails): Presidio usa en_core_web_sm cacheado …` (singleton + `NlpEngineProvider`); (2) el recognizer de `PERSON` de `en_core_web_sm` daba falsos positivos sobre los textos de los escenarios → `01_sip`/`02_mos` salían `blocked_by_guardrail` → fix `fix(policies): detect_pii/pii_leakage solo con recognizers de patrón fiables` (se quitan `PERSON`/`PHONE_NUMBER`); (3) `GET /executions` solo leía el JSONL (sin las ejecuciones `awaiting_approval`) → la página de Aprobaciones nunca veía pendientes → fix `fix(api): /executions incluye también las ejecuciones HITL pendientes` (+ test). Tras los fixes, el walkthrough pasa 26/26: las 5 páginas sirven 200, los 3 escenarios del demo funcionan (MOS→completed, SIP→awaiting→approve→completed / reject→failed), bloqueo PII (NIF/email) OK, historial/violaciones/políticas OK, y el estado sobrevive a `down`+`up`.
|
||||
|
||||
### Task 36: Dockerfiles para core y dashboard
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -4,7 +4,11 @@ description: Política base aplicada a agentes de operación de plataforma de vo
|
||||
input_validators:
|
||||
- type: detect_pii
|
||||
config:
|
||||
entities: [PERSON, EMAIL_ADDRESS, PHONE_NUMBER, ES_NIF, IP_ADDRESS, IBAN_CODE]
|
||||
# Solo recognizers de patrón fiables: el modelo spaCy en_core_web_sm que usa
|
||||
# Presidio en la imagen Docker da falsos positivos de PERSON sobre los textos
|
||||
# técnicos de los escenarios (que están en español y se analizan con el modelo
|
||||
# en). PERSON/PHONE_NUMBER necesitarían en_core_web_lg o un modelo es_*.
|
||||
entities: [EMAIL_ADDRESS, ES_NIF, IP_ADDRESS, IBAN_CODE]
|
||||
severity_on_match: block
|
||||
- type: prompt_injection
|
||||
config:
|
||||
@@ -45,7 +49,7 @@ output_validators:
|
||||
requires_approval: {type: boolean}
|
||||
- type: pii_leakage
|
||||
config:
|
||||
entities: [EMAIL_ADDRESS, PHONE_NUMBER, ES_NIF, IP_ADDRESS]
|
||||
entities: [EMAIL_ADDRESS, ES_NIF, IP_ADDRESS]
|
||||
severity_on_match: block
|
||||
- type: forbidden_action_keywords
|
||||
config:
|
||||
|
||||
@@ -81,6 +81,18 @@ def test_list_executions_incluye_la_completada(client: TestClient) -> None:
|
||||
assert rows[0]["n_proposed_actions"] == 1
|
||||
|
||||
|
||||
def test_list_executions_incluye_la_awaiting(client: TestClient) -> None:
|
||||
"""Una ejecución pausada en HITL no se escribe en el JSONL; aún así debe listarse
|
||||
(la página de Aprobaciones del dashboard depende de ello)."""
|
||||
trace_id = client.post(
|
||||
"/agents/incident_analyzer/invoke", json={"input": "caída registros sip"}
|
||||
).json()["trace_id"]
|
||||
rows = client.get("/executions").json()
|
||||
awaiting = [r for r in rows if r["trace_id"] == trace_id]
|
||||
assert len(awaiting) == 1
|
||||
assert awaiting[0]["status"] == "awaiting_approval"
|
||||
|
||||
|
||||
def test_approve_resume_completa_ejecucion(client: TestClient) -> None:
|
||||
invoked = client.post(
|
||||
"/agents/incident_analyzer/invoke", json={"input": "caída registros sip"}
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
"""Tests del CoreClient: cómo mapea respuestas HTTP del core a dicts.
|
||||
|
||||
Importante para el dashboard: una ejecución correcta (``AgentExecution`` serializado)
|
||||
incluye SIEMPRE el campo ``error`` (``str | None``). El wrapper que el cliente añade
|
||||
ante un 4xx debe usar una clave distinta (``api_error``) para no colisionar con ese
|
||||
campo; si usara ``error``, las páginas no podrían distinguir "el agente terminó sin
|
||||
error" de "la API devolvió 404/409/422".
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Callable
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from agentforge_dashboard.client import CoreClient
|
||||
|
||||
|
||||
def _client_with(handler: Callable[[httpx.Request], httpx.Response]) -> CoreClient:
|
||||
client = CoreClient(base_url="http://test")
|
||||
client._client.close()
|
||||
client._client = httpx.Client(base_url="http://test", transport=httpx.MockTransport(handler))
|
||||
return client
|
||||
|
||||
|
||||
def test_post_2xx_devuelve_el_cuerpo_verbatim_incluido_error_none() -> None:
|
||||
"""Un AgentExecution con ``error=None`` se devuelve tal cual, sin envolver."""
|
||||
execution = {
|
||||
"trace_id": "11111111-1111-1111-1111-111111111111",
|
||||
"status": "completed",
|
||||
"error": None,
|
||||
"violations": [],
|
||||
"decision_path": [],
|
||||
}
|
||||
|
||||
def handler(_: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(200, json=execution)
|
||||
|
||||
result = _client_with(handler).invoke_agent("incident_analyzer", {"input": "x"})
|
||||
assert result == execution
|
||||
assert "api_error" not in result
|
||||
|
||||
|
||||
@pytest.mark.parametrize("status", [404, 409, 422])
|
||||
def test_post_4xx_envuelve_el_detalle_en_api_error(status: int) -> None:
|
||||
def handler(_: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(status, json={"detail": "boom"})
|
||||
|
||||
result = _client_with(handler).invoke_agent("incident_analyzer", {"input": "y"})
|
||||
assert result == {"api_error": {"detail": "boom"}}
|
||||
|
||||
|
||||
def test_post_5xx_propaga_la_excepcion() -> None:
|
||||
def handler(_: httpx.Request) -> httpx.Response:
|
||||
return httpx.Response(500, json={"detail": "kaput"})
|
||||
|
||||
with pytest.raises(httpx.HTTPStatusError):
|
||||
_client_with(handler).invoke_agent("incident_analyzer", {"input": "z"})
|
||||
Reference in New Issue
Block a user