# bus-engine

El bus de nicetry: **la única arista entre engines**. Ningún engine guarda la URL ni la key de otro; lo único que conoce es el bus (`BUS_URL` y su key por tenant). Resuelve dos cosas: **eventos** ("pasó esto", de uno a muchos, asíncrono, con entrega al menos una vez, en orden por subject, reintentos, dead letter y replay) y **comandos y lecturas** ("hacé esto", "dame esto": un registro de servicios y un JWT de 5 minutos firmado por el bus, con el que un engine llama directo a otro). Multi-tenant, REST + MCP. Spec completa: `/openapi.json` (Swagger UI en `/docs`). Cumple el **contrato v1** de engines de nicetry. Diseño: `nicetry/.development/about/bus-engine.md`.

## Auth

- Rutas de tenant: `/{tenant}/...` con el header `X-Engine-Key`. Hay dos clases de key, las dos van en el mismo header:
  - la **key del tenant** (`POST /admin/tenants`): para interfaces, el admin y pruebas; ve y opera todo lo del tenant;
  - la **key de un engine en el tenant** (`POST /admin/tenants/{slug}/engines/{engine}/key`): la que cada engine guarda junto con `BUS_URL`. Publica como ese engine (`source` fijo), maneja solo su propia suscripción y es la única que puede hacer `resolve`.
- Sin una key válida, 401 (exista o no el tenant); con la key de otro tenant, 403.
- Admin: `/admin/...` con `x-admin-api-key`. Las keys se devuelven una sola vez. `DELETE /admin/tenants/{slug}` da 409 si el tenant tiene eventos, salvo con el cuerpo `{"confirm": "<slug>"}`.
- `GET /.well-known/jwks.json` es público: las claves con las que los engines verifican los tokens de resolve.

## Convenciones

- **`X-Actor: <system>:<id>`** en toda escritura (`panel:ana@empresa.com`, `agent:<rutina>`, `loops:<slug>`, `bus:<engine>`). Si falta, vale el `actor` del cuerpo JSON, y si tampoco está, `unknown:`. Las lecturas nunca lo exigen.
- **`X-Request-Id`**: se respeta el que llega o se genera uno; vuelve siempre en la respuesta. Un evento publicado sin `requestId` lleva el del request, y la entrega lo propaga.
- **Errores**: `{ "error": "...", "code": "...", "details": [...], "requestId": "..." }`. `code` es estable: `validation_error`, `invalid_json`, `unauthorized`, `forbidden`, `not_found`, `conflict`, `idempotency_conflict`, `payload_too_large`, `unsupported_media_type`, `rate_limited` (con `Retry-After`), `unavailable` (la base), `internal`.
- **Listados** (`/events`, `/deliveries`): `?limit` (100 por defecto, hasta 500) y `?cursor`. El cuerpo es el array; el cursor de la página siguiente viene en el header `X-Next-Cursor`, ausente en la última.
- **`Idempotency-Key`** en los POST: con la misma key y el mismo cuerpo dentro de 24 h devuelve la respuesta original (`Idempotent-Replayed: true`); con otro cuerpo, 409 `idempotency_conflict`. Para `POST /events` alcanza con el `id` del evento: un id repetido da 202 sin duplicar.
- Cuerpos de hasta 1 MB (413 `payload_too_large`).

## Eventos

**Envelope** (lo que se publica y lo que llega en la entrega):

```json
{
  "id": "01J…",                   // ULID: único global, ordenable, clave de idempotencia (lo genera el emisor o el bus)
  "tenant": "toprentals",
  "type": "hub.card.moved",       // <engine>.<entidad>.<verbo>
  "version": 1,                   // versión del esquema de este type
  "time": "2026-10-07T14:03:11Z", // cuándo pasó en el origen
  "source": "hub",                // engine que publica
  "subject": "card:9a4…",         // entidad principal: clave de orden
  "actor": "panel:ana@empresa.com",
  "requestId": "…",
  "data": { "entityType": "card", "entityId": "9a4…", "fromStage": "…", "toStage": "…" }
}
```

`data` lleva ids y lo mínimo para decidir; contenido solo cuando el consumidor lo necesita para actuar (el texto de un mensaje de WhatsApp), y **nunca secretos**. Quien necesita el estado lo lee del origen (con `resolve`).

- **Publicar**: `POST /{tenant}/events` con un evento o un array de hasta 100. 202 cuando quedó persistido y encolado para cada suscripción que lo matchea; la respuesta trae `ids`, `accepted` y `duplicates`. **Outbox en el emisor**: el engine escribe el evento (con su `id`) en una tabla `outbox` en la misma transacción que el cambio y un proceso lo publica y lo marca; si repite, el bus no duplica.
- **Leer**: `GET /{tenant}/events?type=hub.card.*&subject=&source=&since=&until=` y `GET /{tenant}/events/{id}` (con el estado de sus entregas). Retención: 30 días.
- **Catálogo de types** (`GET /{tenant}/event-types?owner=`): cada type tiene dueño (el único engine que puede publicarlo), `version`, `content` (si `data` lleva contenido y no solo ids), la forma del `subject` y un JSON Schema de `data`. Publicar un type que no está en el catálogo da 400 `unknown_type`; uno de otro dueño o con `data` que no cumple el esquema, 400 `validation_error` con `details` por índice; el lote se acepta o se rechaza entero. Un cambio que rompe es un type o una versión nueva, nunca un cambio en caliente. El admin lo mantiene con `PUT|DELETE /admin/event-types/{type}`. Hoy: `hub.<card|entry|contact|company|appointment|card-link>.<created|updated|moved|deleted|completed|cancelled|no-show>` (ids: `{entityType, entityId}`, subject `<entityType>:<entityId>`) (los types de hub van con guion, `hub.card-link.created`, `hub.appointment.no-show`, aunque la entrega directa de hub siga diciendo `card_link.created` y `appointment.no_show`: lo que viaja igual por los dos caminos es el `id`) y `whatsapp.message.received` (con contenido: `text` o `caption`) y `whatsapp.message.status` (subject `contact:<numberId>:<waId>`), `whatsapp.metadata.updated` (subject `<entityType>:<entityId>`).

## Suscripciones

- `PUT /{tenant}/subscriptions/{consumer}` con `{"types": ["hub.card.*", "whatsapp.message.received"], "filter": {…}, "active": true}`. `consumer` es un engine registrado; `*` cubre un segmento o, al final, el resto; `filter` es igualdad sobre campos de `data`. La **URL de entrega no la elige el suscriptor**: es siempre `{baseUrl registrada del consumer}/{tenant}/inbound`. Al crearla devuelve `secret` (el HMAC de la firma) una sola vez; `POST …/rotate-secret` da uno nuevo. `active: false` pausa sin borrar.
- `GET /{tenant}/subscriptions`, `GET|DELETE /{tenant}/subscriptions/{consumer}` (409 con entregas, salvo `{"confirm": "<consumer>"}`).
- Lo específico de una suscripción (la rutina que despierta un aviso en agents) vive en el consumidor, no acá.

## Entrega

`POST {baseUrl}/{tenant}/inbound` con el envelope en el cuerpo (JSON) y los headers:

- `X-Bus-Event-Id` y `X-Event-Id` (el `id`), `X-Event-Timestamp` (el `time`), `X-Bus-Delivery-Attempt` (1, 2, …), `X-Request-Id`, `X-Actor: bus:<source>`;
- `X-Bus-Signature: t=<unix>,v1=<hmac-sha256(secret, t + "." + cuerpo crudo)>`. El receptor rechaza `t` de más de 5 minutos, compara en tiempo constante y **deduplica por `id`**.

El consumidor responde 2xx cuando procesó o encoló. Cualquier otra respuesta, un 3xx o un timeout de 10 s cuenta como fallo.

- **Reintentos**: 10 s, 1 min, 5 min, 30 min, 2 h (con jitter); 6 intentos en total y después **dead letter** (`status: dead`).
- **Orden**: por `(tenant, consumer, subject)`. Dos eventos del mismo subject llegan en el orden en que se publicaron; eventos de subjects distintos pueden ir en paralelo. Un evento que está fallando bloquea solo a los siguientes de su subject. **Uno muerto no bloquea**: el dead letter destraba el subject, porque los eventos llevan ids y el consumidor lee el estado real del origen, así que un evento muerto no tiene por qué frenar días a los que vienen atrás. Un redeliver de un muerto sale con su id original y el consumidor deduplica como siempre.
- **Garantiza** al menos una vez a cada suscripción activa, en orden por subject, o dead letter visible. **No garantiza** exactamente una vez (el consumidor deduplica por `id`), orden global entre subjects ni tiempo real duro (objetivo: menos de 2 s en condiciones normales).

### Operar las entregas

- `GET /{tenant}/deliveries?status=dead|pending|retrying|delivering|delivered&consumer=&eventId=`: el dead letter es `status=dead`. `GET /{tenant}/deliveries/summary?consumer=`: conteo por estado. `GET /{tenant}/deliveries/{id}`.
- `POST /{tenant}/deliveries/{id}/redeliver`: una muerta vuelve a la cola desde cero; una pendiente o en reintento se intenta ya.
- `POST /{tenant}/subscriptions/{consumer}/replay` con `{"from": "<id o fecha ISO>", "to"?: …, "types"?: [...]}`: vuelve a entregar desde ese punto (backfill, recuperación, alta de un consumidor con historia). Las entregas de replay vienen marcadas `replay: true`.

## Registro de servicios, grants y resolve

- **Registro** (admin, una vez por engine): `PUT /admin/engines/{name}` con `{"baseUrl": "https://…", "scopes": ["roles:read", "messages:send"]}`. `GET /admin/engines`, `DELETE /admin/engines/{name}`.
- **Grants** (admin, por tenant): `PUT /admin/tenants/{slug}/grants/{caller}/{target}` con `{"scopes": [...]}` (subconjunto de los que `target` expone). `GET /admin/tenants/{slug}/grants`, `DELETE …`.
- **Resolve** (solo con la key de un engine): `POST /{tenant}/resolve` con `{"target": "ops", "scopes": ["roles:read"]}` → `{"baseUrl", "token", "expiresAt", "scopes"}`. El token es un JWT **ES256** con `iss=bus`, `sub=<engine que llama>`, `aud=<target>`, `tenant`, `scopes`, `jti`, 5 minutos. El que llama cachea `baseUrl` y `token` hasta `expiresAt` y llama directo: `GET {baseUrl}/{tenant}/roles` con `Authorization: Bearer <token>`.
- **Verificar** (el engine destino): con la JWKS pública `GET {BUS_URL}/.well-known/jwks.json` (cacheada; una clave rotada sigue un día más), chequear firma, `iss=bus`, `aud=<yo>`, `exp`, que `tenant` coincida con el de la ruta y que los `scopes` cubran la operación. Cada engine acepta así dos credenciales: su `X-Engine-Key` del tenant (interfaces y consumidores externos) y el Bearer del bus (otros engines).
- `GET /{tenant}/grants`: lo que este principal puede resolver. `GET /{tenant}/whoami` → `{tenant, engine}` (`engine` null = key del tenant): probe para saber qué key es.
- Rotación de la clave de firma: `POST /admin/signing-keys/rotate`, `GET /admin/signing-keys`.

## Alta de un tenant (lo que hace el admin)

En un paso, idempotente: `POST /admin/tenants/{slug}/provision` con `{"engines": ["hub", "whatsapp", "agents", "loops"], "grants": [{"caller", "target", "scopes"}], "subscriptions": [{"consumer": "agents", "types": ["whatsapp.message.*", "hub.card.*"]}], "rotate": ["tenant", "hub", "agents:secret"]}`. Crea lo que falte (tenant, key de cada engine, grants, suscripciones) y devuelve cada credencial **solo si se creó en esta llamada o se pidió con `rotate`** (`tenantApiKey`, `engines[].apiKey`, `subscriptions[].secret`; el resto null). Repetirlo no cambia nada. Los engines tienen que estar registrados. O por pasos sueltos:
1. `POST /admin/tenants {"slug"}` → key del tenant (bóveda del admin).
2. Para cada engine del tenant: `POST /admin/tenants/{slug}/engines/{engine}/key` → la key que ese engine guarda con `BUS_URL` (rota si ya había); y sus grants.
3. Las suscripciones por defecto con `PUT /{tenant}/subscriptions/{consumer}` → el `secret` va al engine consumidor.

## MCP

`/{tenant}/mcp` (Streamable HTTP, stateless, mismo `X-Engine-Key`). `help` devuelve este documento y `api` llama a cualquier ruta REST del tenant (`method`, `path` relativo al tenant, `query`, `body`), con la misma key y el mismo `X-Actor`. Devuelve `<status>`, después `next-cursor: <c>` si hay más (se sigue con `query.cursor`) y `deprecation: …` si la ruta está deprecada, y al final el cuerpo.
