CRM Adapter — HubSpot

Adapter independiente que recibe tareas del Core Service Bus por Cloud Tasks, hace upsert de contactos en HubSpot CRM y publica el resultado como CloudEvent en Pub/Sub. Nunca se expone a internet.

NestJS 11 · Node 20 HubSpot CRM API v3 Cloud Run · puerto 8080 Entrada: Cloud Tasks Salida: Pub/Sub env: — v—

Resumen

Qué hace

Sincroniza contactos de los sistemas internos de Bebbia hacia HubSpot CRM. Recibe un CloudEvent envuelto en una tarea de Cloud Tasks, lo transforma a propiedades de HubSpot, crea el contacto (o lo actualiza si ya existe) y publica un evento de resultado que el CSB consume para continuar el flujo.

Qué no hace

No expone endpoints públicos, no recibe webhooks de HubSpot, no gestiona deals ni Marketing Hub. Los endpoints de deals y resolución de contactos referenciados por la saga order-to-hubspot siguen pendientes (ver pendientes).

Identidad del servicio

Componente
csb-crm-adapter
Cloud Run
rtp-transversal-dev-dev-crm-adapter (patrón <project>-<env>-crm-adapter)
Imagen
us-central1-docker.pkg.dev/<project>/csb-adapters/crm-adapter:<sha>
Región
us-central1
Proyecto
rtp-transversal-dev

Equipo y origen

Stakeholder
Hugo Quintero (ArchOps)
Desarrollo
Héctor Cruz, Carlos Martínez
Origen
Migrado desde rtp-csb/apps/hubspot-adapter (ago 2026)
Rama activa
dev → despliega con Cloud Build

Arquitectura

El adapter es una pieza del patrón orquestación del CSB: el bus decide cuándo sincronizar, encola la tarea y reacciona al resultado. El adapter solo sabe hablar con HubSpot.

Cómo encaja con el CSB

  1. Un workflow del CSB (saga order-to-hubspot) crea una tarea en la cola de Cloud Tasks del adapter con el SyncPayloadDto como cuerpo.
  2. Cloud Tasks invoca POST /sync con un token OIDC firmado por la cuenta de servicio del CSB (csb-run-sa@…), que tiene roles/run.invoker sobre el servicio.
  3. El adapter transforma, llama a HubSpot y publica el resultado en el topic csb-crm-results-<env>.
  4. El CSB consume la suscripción csb-crm-results-<env>-csb-sub y continúa el flujo.
Acceso El servicio tiene ingress INTERNAL_LOAD_BALANCER (política de organización run.allowedIngress) y solo la SA del CSB tiene run.invoker. Desde internet la URL *.run.app responde 404 antes de llegar al contenedor. Esta documentación se sirve por la misma ruta interna que el resto del servicio.

Módulos internos

MóduloResponsabilidadArchivos clave
SyncModuleRecibe la tarea, orquesta transformación → HubSpot → publicación.sync/sync.controller.ts, sync/sync.service.ts
HubSpotTransformerConvierte el envelope CloudEvents en propiedades de contacto HubSpot. Exige email.sync/transformer.ts
HubSpotClientServiceCliente HTTP (fetch nativo) con Bearer token. Estrategia crear → si 409, buscar por email → PATCH.sync/hubspot-client.ts
ResultPublisherServicePublica el CloudEvent hubspot.contact.synced.v1 en Pub/Sub.result-publisher/result-publisher.service.ts
HealthModuleLiveness con ping a HubSpot y Secret Manager; readiness simple.health/health.controller.ts
DocsModuleSirve este portal, la guía de despliegue y registra Swagger UI.docs/docs.controller.ts, docs/public/
config/Carga de configuración dual: variables de entorno en local, Secret Manager en Cloud Run.config/app.config.ts, config/swagger.config.ts
common/Filtro global de excepciones (JSON a stderr) y logger estructurado.common/filters/http-exception.filter.ts, common/logger.service.ts

Bootstrap (main.ts): ValidationPipe global con transform y whitelist (campos desconocidos se descartan), filtro global de excepciones, Swagger bajo /docs/api cuando la documentación está habilitada, escucha en PORT (8080).

Flujo de sincronización

  1. Cloud Tasks → POST /syncCuerpo SyncPayloadDto. Headers inyectados por la cola: Authorization: Bearer <oidc>, X-Task-Queue, X-Retry-Count, X-Workflow-Execution.
  2. ValidaciónValidationPipe valida envelope y taskMetadata. Si falla: 400 y la tarea reintenta (ver errores).
  3. TransformaciónHubSpotTransformer.transform() lee envelope.data y arma properties. Sin email lanza error → 500.
  4. Crear contactoPOST https://api.hubapi.com/crm/v3/objects/contacts con Bearer token, timeout 15 s. Éxito → action: created.
  5. Si HubSpot responde 409 (ya existe)POST /crm/v3/objects/contacts/search filtrando email EQ → toma el primer id → PATCH /crm/v3/objects/contacts/{id}. Éxito → action: updated.
  6. Publicar resultadoResultPublisherService.publishSuccess() publica el CloudEvent en PUBSUB_RESULT_TOPIC. Si la publicación falla, la tarea completa falla (500) y Cloud Tasks reintenta.
  7. Respuesta 200 { "status": "ok" }Cloud Tasks marca la tarea como completada. Todo el proceso deja logs con event_id y duration_ms.
Idempotencia El adapter no deduplica por envelope.id: un reintento vuelve a ejecutar el upsert. Es seguro porque HubSpot devuelve 409 en el segundo intento y el flujo termina en updated, pero se publica un segundo evento de resultado con el mismo id (<originalEventId>-result). El consumidor del CSB debe tolerarlo.

Endpoints

Todos responden JSON. La referencia interactiva con esquemas generados desde los DTOs está en Swagger UI; el documento OpenAPI 3.1 en /docs/openapi.json.

POST/syncInvocado por Cloud Tasks (CSB)

Sincroniza un contacto a HubSpot a partir de un CloudEvent. Responde 200 solo cuando el contacto quedó creado/actualizado y el resultado fue publicado.

Request body — SyncPayloadDto
{
  "envelope": {
    "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890",   // UUID v4 del evento original
    "source": "//bebbia/orders-service",
    "type": "com.bebbia.order.placed.v1",
    "specversion": "1.0",
    "entityid": "ORD-000123",                          // ordering key
    "datacontenttype": "application/json",               // opcional
    "time": "2026-09-09T15:04:05.000Z",                  // opcional
    "data": {
      "email": "ana.lopez@example.com",               // requerido por el transformer
      "first_name": "Ana",
      "last_name": "López",
      "phone": "+52 55 1234 5678",
      "company": "Bebbia"
    }
  },
  "taskMetadata": {
    "retryCount": 0,                                   // requerido, >= 0
    "queueName": "csb-hubspot-adapter-queue-dev",      // opcional
    "taskName": "projects/.../tasks/..."              // opcional
  }
}
Response 200
{ "status": "ok" }
Códigos de respuesta
200 sincronizado y publicado 400 payload inválido (ValidationPipe) 500 falla HubSpot / transformación / Pub/Sub
Ejemplo local
curl -s -X POST http://localhost:8080/sync \
  -H "Content-Type: application/json" \
  -d '{"envelope":{"id":"11111111-1111-4111-8111-111111111111","source":"//bebbia/orders-service","type":"com.bebbia.order.placed.v1","specversion":"1.0","entityid":"ORD-1","data":{"email":"ana@example.com","first_name":"Ana"}},"taskMetadata":{"retryCount":0}}'
GET/healthLiveness probe de Cloud Run · Cloud Monitoring

Ejecuta en paralelo un ping a HubSpot (GET /crm/v3/objects/contacts?limit=1, timeout 5 s) y, si USE_SECRET_MANAGER=true, una lectura del secreto del token. Siempre responde 200; el campo status indica ok o degraded.

{
  "status": "ok",
  "timestamp": "2026-09-09T15:04:05.000Z",
  "version": "0.1.0",
  "component": "csb-crm-adapter",
  "checks": {
    "hubspot_connectivity": "ok",
    "secret_manager": "ok"
  }
}
// status: ok | degraded — siempre HTTP 200
GET/health/readyStartup probe de Cloud Run

Readiness sin dependencias externas. Cloud Run solo marca una revisión como Ready cuando esta ruta responde 200; el pipeline usa ese estado como smoke test.

{ "status": "ok", "timestamp": "2026-09-09T15:04:05.000Z" }
GET/docs · /docs/api · /docs/openapi.json · /docs/how-it-worksPersonas

Portal de documentación (esta página), Swagger UI, documento OpenAPI y guía de integración/despliegue. GET / redirige a /docs. Activos fuera de production; DOCS_ENABLED=true|false fuerza el estado. Deshabilitados responden 404.

Contratos y payloads

Entrada — SyncPayloadDto (Cloud Tasks)

CampoTipoReq.Descripción
envelope.idstringsíUUID v4 del evento de dominio. Se propaga a HubSpot como csb_event_id y al evento de resultado.
envelope.sourcestringsíURI del servicio emisor, p. ej. //bebbia/orders-service. Se guarda como csb_source.
envelope.typestringsíTipo CloudEvents del evento original.
envelope.specversionstringsíSiempre "1.0".
envelope.entityidstringsíIdentificador de la entidad de negocio (ordering key).
envelope.dataobjectsíPayload de negocio. El transformer lee email (obligatorio), first_name, last_name, phone, company.
envelope.datacontenttypestringnoDefault documentado application/avro; en la práctica el CSB envía JSON.
envelope.timedate-timenoTimestamp del evento original.
taskMetadata.retryCountnumber ≥ 0síIntentos previos de entrega. Solo se usa para logging.
taskMetadata.queueNamestringnoNombre corto de la cola.
taskMetadata.taskNamestringnoNombre completo de la tarea en Cloud Tasks.

Mapeo a HubSpot — propiedades del contacto

Propiedad HubSpotOrigenNotas
emaildata.emailObligatoria. Clave de deduplicación en HubSpot (409 si ya existe).
firstnamedata.first_nameCadena vacía si no viene.
lastnamedata.last_nameCadena vacía si no viene.
phonedata.phoneSin normalización.
companydata.companyTexto libre, no crea objeto Company.
csb_event_idenvelope.idPropiedad custom: debe existir en el portal de HubSpot.
csb_sourceenvelope.sourcePropiedad custom.
csb_synced_atahora (ISO 8601)Propiedad custom.
Prerequisito en HubSpot Si las propiedades csb_event_id, csb_source y csb_synced_at no existen en el portal, HubSpot responde 400 "Property values were not valid" y todas las tareas fallan. La guía de despliegue detalla cómo crearlas.

Salida — CloudEvent de resultado (Pub/Sub)

Publicado en el topic PUBSUB_RESULT_TOPIC (Terraform: csb-crm-results-<env>) como mensaje JSON con atributos ce_type, ce_source, ce_specversion.

{
  "specversion": "1.0",
  "type": "hubspot.contact.synced.v1",
  "source": "//csb-crm-adapter/rtp-transversal-dev",
  "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890-result",   // <originalEventId>-result
  "time": "2026-09-09T15:04:06.120Z",
  "datacontenttype": "application/json",
  "data": {
    "originalEventId": "a1b2c3d4-e5f6-7890-abcd-ef1234567890",
    "email": "ana.lopez@example.com",
    "hubspotId": "901234567",
    "action": "created",                                  // created | updated
    "status": "SUCCESS"
  }
}

Solo se publican éxitos. Un fallo no genera evento: la tarea responde 500 y Cloud Tasks reintenta; agotados los reintentos, la cola descarta la tarea (no hay DLQ configurada en este adapter).

Errores y reintentos

Formato de error (HttpExceptionFilter)

{
  "statusCode": 500,
  "timestamp": "2026-09-09T15:04:05.000Z",
  "path": "/sync",
  "error": { "message": "Internal server error" }   // en 400: objeto de ValidationPipe con "message": [...]
}

Qué produce cada código

  • 400Envelope o taskMetadata inválidos (campos faltantes, tipos incorrectos). Cloud Tasks reintenta igual: corregir el productor.
  • 500Missing required field: email; HubSpot no-2xx distinto de 409 (auth 401, rate limit 429, 400 por propiedades inexistentes); contacto no encontrado tras 409; timeout 15 s; fallo al publicar en Pub/Sub. El mensaje real queda en logs, no en la respuesta.
  • 200Contacto creado o actualizado y evento publicado.

Política de reintentos (Cloud Tasks, Terraform)

max_attempts
5
backoff
10 s → 300 s, max_doublings 4
max_retry_duration
3600 s
rate
10 dispatches/s · 5 concurrentes

Códigos de error declarados en ErrorResponseDto (HUBSPOT_TIMEOUT, HUBSPOT_AUTH_FAILED, HUBSPOT_RATE_LIMITED, SCHEMA_VALIDATION_FAILED) describen la intención del contrato; hoy el filtro global no los emite en la respuesta y el log usa error_code: HUBSPOT_SYNC_ERROR.

Integración con el CSB: cómo lo consume el bus

De extremo a extremo: qué sistema dispara, por qué puerta entra al Core Service Bus, cómo el bus localiza e invoca a este adapter, qué flujos lo usan hoy, qué hace con el resultado y qué le falta a cada pieza. Fuente: rtp-csb en dev (seed, workflows/, infra/envs/dev) al 10 sep 2026.

Mapa de extremo a extremo

Regla de la arquitectura: solo el CSB le pega a HubSpot. Ningún sistema externo llama al adapter; publican un evento o llaman al WebHook del CSB, y es un Cloud Workflow (saga) quien decide cuándo y con qué invocar a este adapter.

Quién dispara y por qué puerta entra al CSB

El CSB tiene tres puertas de entrada. Las dos primeras son para sistemas; la tercera para personas. Un externo (tienda Bebbia, Data Mesh, Portal de Inventario, CommerceTools) nunca conoce la URL del adapter: conoce un topic o un WebHook del CSB.

PuertaQuién la usaAutenticaciónCuándoLlega a este adapter vía
Topic Pub/Sub con schema
order.placed.v1, order.cancelled.v1, order.fulfillment.requested.v1
Servicio de órdenes / tienda (productor interno con SA)IAM roles/pubsub.publisher sobre el topic; el mensaje se valida contra el schema Avro/JSON registradoAl crear, cancelar o solicitar el cumplimiento de una ordenEventarc dispara el Cloud Workflow correspondiente; el workflow llama al adapter
WebHook de flujo
POST https://csb-dev.rotoplas.com/api/flows/{flowId}/webhook
Sistemas sin cuenta GCP: tienda Bebbia (CSB-02), Bringoz (CSB-01)HMAC-SHA256 del cuerpo con el secreto csb-<env>-webhooks-signing-key, header X-CSB-SignatureCuando el externo tiene un hecho que reportar y quiere respuesta síncronaHoy ningún flujo con WebHook usa este adapter; sería igual: el workflow del flujo llama a /sync
Ejecución manual / editor visual
POST …/api/provisioning/flows/{id}/execute
Personas con JWT del CSB (editor en csb-dev.rotoplas.com/flows)Login CSB (Google OAuth / usuario seed)Pruebas, demos, reproceso puntualFlujo order-fulfillment-demo llama a CRM_ADAPTER_URL/sync tres veces

Ejemplo 1 · publicar el evento que dispara order-to-hubspot-saga

Contrato Avro order.placed.v1 (schemas/order.placed.v1.avsc). El envelope CloudEvents lo agrega el CSB; el productor solo publica el payload de dominio.

{
  "event_id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890",
  "order_id": "ORD-000123",
  "customer_id": "CUST-000045",
  "items": [
    {
      "sku_code": "BEB-PURIF-01",
      "product_name": "Purificador Bebbia",
      "quantity": 1,
      "unit_price": "1499.00",
      "total_price": "1499.00"
    }
  ],
  "total_amount": "1499.00",
  "currency": "MXN",
  "placed_at": 1788000000000,
  "channel": "WEB",
  "correlation_id": "corr-7f3a"
}
gcloud pubsub topics publish order.placed.v1 \
  --project=rtp-transversal-dev \
  --ordering-key=ORD-000123 \
  --message="$(cat order-placed.json)"

Ejemplo 2 · llamar un WebHook de flujo (patrón de la tienda Bebbia)

BODY='{"requestId":"REQ-0001","susid":"SUS-100045","storeId":"bebbia-web", ...}'
SIG=$(printf '%s' "$BODY" | openssl dgst -sha256 -hmac "$WEBHOOKS_SIGNING_KEY" | awk '{print $2}')

curl -s -X POST "https://csb-dev.rotoplas.com/api/flows/$FLOW_ID/webhook" \
  -H "Content-Type: application/json" \
  -H "X-CSB-Signature: $SIG" \
  -d "$BODY"
# executeFlow es síncrono: la respuesta ya trae el resultado del workflow

Cómo llega el CSB a este adapter (cableado real)

  1. El adapter publica sus datos de contacto en su propio TerraformOutputs en rtp-crm-adapter/infra/envs/dev: service_url, queue_name (csb-hubspot-adapter-queue-dev), result_topic_id (csb-crm-results-dev). Estado en gs://rtp-transversal-dev-tfstate/crm-adapter/dev.
  2. El CSB los lee por remote state, no los hardcodeartp-csb/infra/envs/dev/main.tf (commit 40352ce): data "terraform_remote_state" "crm_adapter" → local.adapters.crm = { url, queue, tasks_sa = null, result_topic }. Output adapter_endpoints muestra lo resuelto.
  3. La URL entra a csb-api y al job de seed como variable de entornomodule.compute.extra_env.CRM_ADAPTER_URL se inyecta en el Cloud Run csb-api y en el job csb-api-migrate.
  4. El seed registra el adapter en el registry del CSBapps/api/src/database/seed.ts: ApiDefinition HubSpot CRM v3 API (baseUrl = CRM_ADAPTER_URL), Connector hubspot-order-connector (mapping order_id→dealname, total_amount→amount, customer_id→associatedContactId; throttle 10/s, 15 concurrentes, 3 intentos) y Adapter csb-hubspot-adapter con cloudRunUrl = CRM_ADAPTER_URL, probe GET /health. Visible en csb-dev.rotoplas.com/registry/adapters.
  5. Identidad: quién puede invocar al adapterEste adapter usa el modelo inverso: su Terraform otorga roles/run.invoker a csb-run-sa@rtp-transversal-dev.iam.gserviceaccount.com. Cloud Workflows y Cloud Tasks del CSB deben firmar el OIDC con esa SA; si el workflow corre con otra SA, la llamada recibe 403.
  6. Los flujos referencian el adapter por id y el workflow lo llamaEn el grafo del flujo el nodo adapter apunta a csb-hubspot-adapter (refId). El YAML del workflow lee HUBSPOT_ADAPTER_URL con sys.get_env y hace http.post con auth: { type: OIDC } y header X-Workflow-Execution. Los pasos del editor con auth: 'OIDC' hacen lo mismo.
  7. Cloud Tasks cuando el trabajo es diferidoContrato HUBSPOT_ADAPTER_QUEUE en packages/shared/src/contracts/cloud-tasks.contracts.ts: cola csb-hubspot-queue (Terraform del CSB, 100/s, 10 concurrentes, 5 intentos), destino POST {HUBSPOT_ADAPTER_URL}/sync, cuerpo SyncPayloadDto, headers Authorization: Bearer <oidc>, X-Task-Queue, X-Retry-Count, X-Workflow-Execution. deployFlow crea además una cola que-<flow>-<adapter> por cada edge enqueue del editor.
  8. El resultado vuelve por Pub/SubEl adapter publica hubspot.contact.synced.v1 en csb-crm-results-dev; su Terraform crea la suscripción csb-crm-results-dev-csb-sub para que el CSB la consuma.
Cómo probar que el cableado está vivo terraform output -json adapter_endpoints en rtp-csb/infra/envs/dev debe mostrar la URL del adapter; en el registry, POST /api/adapters/{id}/probe hace GET <cloudRunUrl>/health con header x-csb-probe: true y guarda lastHealthStatus.

Flujos del CSB que consumen este adapter

Flujo (seed / Terraform)DisparadorMomento de negocioLlama aEstado real
order-to-hubspot-saga
workflow Terraform + flow ACTIVE
order.placed.v1 (Eventarc)Se creó una orden: registrar al cliente y la venta como deal en HubSpotPOST /contacts/resolve → POST /deals/transform → POST /deals/syncfalla las tres rutas no existen en el adapter
order-cancellation-notify
workflow Terraform + flow ACTIVE
order.cancelled.v1Se canceló una orden: cerrar el dealPOST /deals/sync con envelope <id>-cancelfalla ruta inexistente
order-fulfillment-orchestrator
workflow Terraform
order.fulfillment.requested.v1Cumplimiento completo: reserva SAP + deal HubSpot con compensación/contacts/resolve, /deals/sync (compensación action: CANCEL_DEAL)falla rutas inexistentes
order-fulfillment-demo
flow del editor (stepGraph)
manual / order.fulfillment.requested.v1Demo end-to-end desde el editor visualPOST /sync ×3 (contacto, reintento, deal)parcial llega a /sync pero el envelope no cumple SyncPayloadDto
Suscripción push order.placed.v1 → hubspot-adapter
Terraform subscription_order_hubspot
order.placed.v1Camino directo sin workflowPOST <adapter interno del CSB>/sync (mensaje Pub/Sub crudo)legado apunta al Cloud Run hubspot-adapter del repo CSB, no a este servicio

Secuencia real de order-to-hubspot-saga (como está escrita)

  1. Eventarc → Cloud WorkflowEl topic order.placed.v1 dispara el workflow con args.envelope y args.taskMetadata.
  2. validate_payloadExige envelope.id, envelope.data.order_id, envelope.data.customer_id; si falta algo lanza VALIDATION_ERROR.
  3. resolve_hubspot_contactPOST {HUBSPOT_ADAPTER_URL}/contacts/resolve con { customerId, email }; espera hubspotContactId.
  4. transform_to_hubspotPOST /deals/transform con { envelope, hubspotContactId, preview: true }; espera data con el deal.
  5. sync_to_hubspotPOST /deals/sync con { envelope, dealPayload, hubspotContactId } y header X-Workflow-Execution; espera dealId.
  6. handle_success / handle_failureÉxito: retorna { status: SUCCESS, orderId, hubspotDealId, hubspotContactId } (no publica en Pub/Sub). Fallo: publica el envelope en order.placed.v1.dlq y retorna FAILED.

Lo que este adapter sí sabe hacer hoy: POST /sync

Cuerpo exacto que espera (el que envía Cloud Tasks según el contrato HUBSPOT_ADAPTER_QUEUE):

{
  "envelope": {
    "id": "a1b2c3d4-e5f6-7890-abcd-ef1234567890",
    "source": "//bebbia/orders-service",
    "type": "com.bebbia.order.placed.v1",
    "specversion": "1.0",
    "entityid": "ORD-000123",
    "data": {
      "email": "ana.lopez@example.com",
      "first_name": "Ana",
      "last_name": "López",
      "phone": "+52 55 1234 5678",
      "company": "Bebbia"
    }
  },
  "taskMetadata": {
    "retryCount": 0,
    "queueName": "csb-hubspot-queue"
  }
}

El flujo order-fulfillment-demo envía en cambio { "envelope": { "id", "entityid", "correlationid", "data": { "customer_id", "email" } } }: sin source, type, specversion ni taskMetadata, por lo que el ValidationPipe responde 400. El paso de deal ni siquiera trae email.

Qué pasa con el resultado

Lo que publica el adapter

CloudEvent hubspot.contact.synced.v1 en csb-crm-results-dev (atributos ce_type, ce_source, ce_specversion) con originalEventId, email, hubspotId, action y status: SUCCESS. Solo éxitos; los fallos se reintentan desde Cloud Tasks.

Lo que el CSB espera

El seed registra el topic hubspot.order.synced.v1 (JSON Schema: event_id, order_id, hubspot_deal_id, hubspot_contact_id, synced_at, …) como salida de order-to-hubspot-saga. Ni el workflow ni el adapter publican ahí.

Quién consume

La suscripción csb-crm-results-dev-csb-sub existe (Terraform del adapter) pero ningún consumer del CSB la lee. Los consumers reales (audit, analytics, notification) están suscritos a inventory.stock.updated.v1.

Trazabilidad

Cada ejecución queda en workflow_executions (GET /api/monitoring/workflows/executions), las tareas en GET /api/monitoring/tasks y los mensajes fallidos en GET /api/monitoring/dlq con replay. Los logs del adapter llevan event_id = envelope.id.

Cómo lo consume un externo (tienda Bebbia, Data Mesh, otro dominio)

Un sistema externo nunca llama al adapter. Tiene dos maneras de usar lo que el adapter hace y una de enterarse del resultado:

Necesidad del externoPatrón en el CSBEjemplo vivo en el repoPasos
Quiero que un hecho de mi sistema termine en HubSpot (fire-and-forget)Publicar en un topic con schema; una saga hace el restoorder.placed.v1 → order-to-hubspot-sagaPedir pubsub.publisher sobre el topic, cumplir el Avro, incluir event_id único y correlation_id
Quiero enviar una petición y recibir la respuesta en la misma llamadaWebHook de flujo firmado (HMAC) → executeFlow síncronoCSB-02 alta-cliente-tienda (tienda → SAP WS1 → callback)Acordar el inputTemplate del flujo, recibir el flowId y el secreto de firma, firmar cada cuerpo
Quiero enterarme cuando HubSpot ya tiene el contacto (Data Mesh, analítica, otro dominio)Suscripción pull propia al topic de resultados, con DLQ e idempotencia por event_idConsumers audit/analytics sobre inventory.stock.updated.v1 (módulo csb-subscription)Terraform csb-subscription (sub + DLQ + SA) sobre csb-crm-results-dev, deserializar el CloudEvent, dedup por id
Quiero que el CSB me avise por HTTP (callback)Etapa de notificación en el flujo (external-api node) con reintentosCSB-02 notificar-alta-tienda → POST /v1/customers/{susid}/signup-resultExponer el endpoint, registrarlo como ApiDefinition con su credencial en /registry/apis, agregar la etapa al flujo

Data Mesh hoy es destino, no productor: CSB-01 le hace POST /v1/domains/inventory/assets/{serialNumber}/status-events (ApiDefinition Core Data Mesh Rotoplas, token BEARER, secreto csb-dev-data-mesh-credential) y usa su lastEventId como candado de idempotencia. Para que consuma resultados de HubSpot aplica la tercera fila.

Cuándo · para qué · por qué

Cuándo (evento)Para qué (resultado de negocio)Por qué pasa por el CSB y no directo
Se coloca una orden (order.placed.v1)Cliente como contacto y venta como deal en HubSpot, con el event_id del CSB en propiedades customReintentos y backoff sin código en la tienda; DLQ y replay; la tienda no guarda tokens de HubSpot; auditoría por event_id
Se cancela una orden (order.cancelled.v1)Cerrar el dealMismo canal y misma trazabilidad; compensación declarada en el saga
Se pide cumplimiento (order.fulfillment.requested.v1)Reserva SAP y deal HubSpot como una sola transacción con compensaciónSolo un orquestador puede deshacer el paso anterior si el siguiente falla
Una persona ejecuta desde el editorValidar un flujo nuevo o reprocesar un casoEl editor versiona el grafo, hace preview de recursos y guarda la ejecución paso a paso

Paso a paso: conectar (o reconectar) este adapter en el CSB

  1. Confirmar outputs del adapter
    cd rtp-crm-adapter/infra/envs/dev && terraform output
    # service_url, queue_name, result_topic_id, crm_adapter_sa
  2. Confirmar que el CSB los resuelve
    cd rtp-csb/infra/envs/dev && terraform output -json adapter_endpoints | jq .crm
  3. Alinear la identidadLa SA con run.invoker en este adapter es csb-run-sa@…; la de csb-api es sa-csb-api-dev@… y la de cada workflow la crea el módulo csb-workflow (gcloud workflows describe order-to-hubspot-saga --format='value(serviceAccount)'). Otorgar run.invoker a las SAs reales en rtp-crm-adapter/infra/envs/dev/terraform.tfvars (o convertir csb_invoker_sa_email en lista).
  4. Inyectar la URL en los workflowsAgregar user_env_vars = { HUBSPOT_ADAPTER_URL = local.adapters.crm.url } al módulo csb-workflow para order-to-hubspot-saga, order-cancellation-notify y order-fulfillment-orchestrator.
  5. Sembrar el registryEl job csb-api-migrate corre el seed con CRM_ADAPTER_URL. En local: CRM_ADAPTER_URL=http://localhost:8082 npm run seed en apps/api. Revisar en /registry/adapters que csb-hubspot-adapter tenga la URL y el probe en verde.
  6. Exponer en el adapter lo que los flujos llamanImplementar POST /contacts/resolve, POST /deals/transform y POST /deals/sync (ya en NEXT_STEPS.md) o cambiar los pasos de los workflows para usar /sync con SyncPayloadDto completo.
  7. Construir o ajustar el flujo en el editor/flows/:id: nodo topic → workflow (pasos call con auth: OIDC hacia {{CRM}}/sync) → nodo adapter (csb-hubspot-adapter) → topic de resultado. Preview GCP muestra la cola que-<flow>-csb-hubspot-adapter; Deploy la crea.
  8. Ejecutar y observarPOST /api/provisioning/flows/{id}/execute con el inputTemplate; ver /monitoring (executions, tasks, dlq) y en el adapter los logs sync_request_received → sync_completed.
  9. Cerrar el ciclo del resultadoDecidir un solo topic de resultado (csb-crm-results-dev con CloudEvent hubspot.contact.synced.v1, o hubspot.order.synced.v1 del seed) y crear el consumer o la suscripción del CSB que lo lea.

Brechas entre lo que el CSB espera y lo que el adapter expone

  • bloqueanteRutas inexistentes. Los tres workflows Terraform llaman a /contacts/resolve, /deals/transform y /deals/sync; el adapter solo tiene /sync. Toda ejecución termina en CONTACT_RESOLVE_ERROR y el envelope va a order.placed.v1.dlq.
  • bloqueanteContrato de /sync. Exige envelope CloudEvents completo (source, type, specversion, entityid) y taskMetadata.retryCount; los pasos del editor (order-fulfillment-demo) envían un envelope mínimo → 400. Además data.email es obligatorio y el paso de deal no lo trae → 500.
  • legadoDos adapters HubSpot en paralelo. csb-compute aún declara el Cloud Run hubspot-adapter interno y la suscripción push order.placed.v1 → hubspot-adapter/sync; los workflows leen HUBSPOT_ADAPTER_URL. Hay que apuntar ambos a local.adapters.crm.url y retirar el servicio interno.
  • deudaTopic de resultado sin acuerdo. Adapter: csb-crm-results-dev / hubspot.contact.synced.v1. Seed: hubspot.order.synced.v1. Nadie consume ninguno.
  • bloqueanteIdentidad. El adapter autoriza a csb-run-sa@…, pero la SA de csb-api es sa-csb-api-dev@… y cada Cloud Workflow corre con su propia SA (módulo csb-workflow). Ninguna de las dos tiene run.invoker aquí: toda llamada OIDC desde el CSB recibe 403.
  • bloqueanteVariables de los workflows. Los YAML leen HUBSPOT_ADAPTER_URL con sys.get_env, pero el módulo csb-workflow no declara user_env_vars: el workflow desplegado no conoce la URL del adapter.
  • deudaPush a /events. Cuando un flujo del editor conecta un topic directo a un nodo adapter, deployFlow crea una suscripción push a <cloudRunUrl>/events, ruta que no existe en ningún adapter.
  • hechoRemote state, CRM_ADAPTER_URL en csb-api y seed, registro en el registry, contrato HUBSPOT_ADAPTER_QUEUE, DLQ y monitoring.

Configuración

loadHubSpotConfig() opera en dos modos: con USE_SECRET_MANAGER=false lee todo de variables de entorno (local, docker compose); con true (Cloud Run) obtiene el token de Secret Manager usando el sufijo de ambiente derivado de NODE_ENV (production → prod).

Variables de entorno

VariableDefaultQuién la defineDescripción
PORT8080DockerfilePuerto HTTP.
NODE_ENVdevelopmentTerraform (dev/qa/production)Ambiente. Deriva el sufijo de secretos y el estado por defecto de la documentación.
USE_SECRET_MANAGERfalseTerraform → trueLeer el token desde Secret Manager.
GOOGLE_CLOUD_PROJECT—TerraformProyecto para Secret Manager.
GCP_PROJECT_IDrtp-transversal-devTerraformProyecto usado por el publicador de Pub/Sub.
HUBSPOT_API_TOKEN—.env localToken de Private App. En Cloud Run viene del secreto.
HUBSPOT_BASE_URLhttps://api.hubapi.comdocker compose (mock)Base URL del cliente; permite apuntar al mock HubSpot local.
PUBSUB_RESULT_TOPICcsb-crm-resultsTerraform → csb-crm-results-<env>Topic donde se publica el resultado. Es el que realmente usa el publicador.
RESULT_TOPIChubspot.contact.synced.v1.env.exampleSe carga en la config pero el publicador no lo usa (ver deuda).
DOCS_ENABLEDautomanualtrue/false fuerza el portal y Swagger. Sin definir: activo salvo en production.

Secretos (Secret Manager)

SecretoContenidoQuién lo crea
csb-<env>-hubspot-api-tokenToken de la Private App de HubSpot (scopes crm.objects.contacts.read/write).Estructura: Terraform (crm-secrets). Valor: manual con gcloud secrets versions add.

La SA del servicio (csb-crm-adapter-<env>@…) tiene secretmanager.secretAccessor solo sobre ese secreto y pubsub.publisher a nivel proyecto.

Infraestructura (Terraform)

Toda la infraestructura vive en infra/ con un módulo de cómputo, uno de secretos y una carpeta por ambiente (envs/dev, envs/qa, envs/prod). Estado remoto en GCS. Cloud Build aplica envs/dev en cada push a dev.

Cloud Run (crm-compute)

Nombre
<project>-<env>-crm-adapter
Ingress
INTERNAL_LOAD_BALANCER
SA
csb-crm-adapter-<env>
Recursos
1 vCPU · 512 Mi · startup CPU boost
Escalado dev
min 0 · max 3
Startup probe
/health/ready cada 5 s, 10 fallos
Liveness
/health cada 30 s, 3 fallos

Cloud Tasks

Cola
csb-hubspot-adapter-queue-<env>
Rate
10/s · 5 concurrentes
Retry
5 intentos · 10 s–300 s · 1 h máx.
Logging
Stackdriver sampling 100 %

Pub/Sub

Topic
csb-crm-results-<env> · retención 7 d
Suscripción
csb-crm-results-<env>-csb-sub (para el CSB) · ack 30 s · retry 10 s–300 s

IAM

Invoker
csb-run-sa@rtp-transversal-dev.iam.gserviceaccount.com → roles/run.invoker
SA adapter
roles/pubsub.publisher (proyecto) · secretAccessor (secreto)

Outputs de Terraform

  • service_url — registrar en el CSB como CRM_ADAPTER_URL (destino de las tareas).
  • queue_name, result_topic_id, crm_adapter_sa.

CI/CD — Cloud Build

Trigger: push a dev en GitLab. Pipeline en cloudbuild.yaml, máquina E2_HIGHCPU_8, timeout 30 min. La imagen se etiqueta con $SHORT_SHA y dev-latest.

  1. security-audit · lint · testEn paralelo. npm audit --audit-level=high, npm run lint, npm run test:cov con umbral de cobertura global 65 % (líneas y funciones).
  2. build · pushDocker multi-stage (deps → builder → runner non-root). Push a Artifact Registry csb-adapters/crm-adapter.
  3. terraform init → plan → applyEn infra/envs/dev con image_tag=$SHORT_SHA y csb_invoker_sa_email.
  4. deploy --no-trafficNueva revisión sin tráfico.
  5. smoke-testEspera (12 × 5 s) a que la revisión quede Ready=True. No hay llamada HTTP: el ingress interno impide que Cloud Build alcance el servicio.
  6. promote-trafficupdate-traffic --to-latest solo si el smoke test pasó.
Documentación en cada deploy Los HTML de src/docs/public/ se copian a dist/docs/public/ en el build (nest-cli.json → compilerOptions.assets) y viajan dentro de la imagen. Un commit a dev que toque la documentación la publica con la siguiente revisión.

Observabilidad

Los logs salen por stdout/stderr y los recoge Cloud Logging. Las excepciones HTTP se escriben como JSON en stderr con severity: ERROR; el resto usa el Logger de NestJS con objetos que incluyen siempre component: "csb-crm-adapter".

Acciones registradas

actionNivelCamposCuándo
sync_request_receivedINFOevent_idEntra la petición al controller.
sync_task_receivedINFOevent_id, retry_count, status: PROCESSINGInicio del procesamiento.
contact_updatedINFOhubspot_idCamino 409 → PATCH.
result_published / result_publish_failedINFO / ERRORtopic, event_id, error_messagePublicación en Pub/Sub.
sync_completedINFOduration_ms, hubspot_id, hubspot_actionÉxito.
sync_failedERRORduration_ms, error_code, error_messageCualquier fallo del flujo.
http_exceptionERRORstatus_code, path, errorFiltro global.

Consulta útil en Cloud Logging

resource.type="cloud_run_revision"
resource.labels.service_name="rtp-transversal-dev-dev-crm-adapter"
(jsonPayload.action="sync_failed" OR severity>=ERROR)

Métricas nativas de Cloud Run (request count por clase de código, latencias, instancias) están disponibles sin configuración extra. Este repo no declara políticas de alerta propias.

Runbook

Las tareas fallan con 500 y sync_failed

  • HubSpot API 401: token inválido o revocado. Subir nueva versión del secreto y redeplegar (la config se lee al arrancar).
  • HubSpot API 400 … Property values were not valid: faltan las propiedades custom csb_* en el portal.
  • HubSpot API 429: rate limit. Bajar max_dispatches_per_second en la cola.
  • Missing required field: email: el productor no envía data.email.

/health reporta degraded

  • hubspot_connectivity: error: HubSpot no alcanzable o token inválido. El liveness sigue en 200, no reinicia el contenedor.
  • secret_manager: error: la SA no tiene secretAccessor o el secreto no tiene versiones.

La revisión nunca queda Ready (pipeline falla en smoke-test)

  • Con USE_SECRET_MANAGER=true el arranque lee el secreto: si no existe o no hay permiso, el contenedor muere. Revisar logs de la revisión.
  • Revisar que PORT sea 8080 y que el startup probe apunte a /health/ready.

Ver la documentación o Swagger desplegados

  • El servicio no es alcanzable desde internet. En dev está publicada por el balanceador del CSB: csb-dev.rotoplas.com/adapters/crm/docs (solo el prefijo /docs; los endpoints de negocio no se exponen).
  • En local: npm run dev y abrir http://localhost:8080/docs.

Reenviar una tarea manualmente

gcloud tasks create-http-task --queue=csb-hubspot-adapter-queue-dev \
  --location=us-central1 --url="$CRM_ADAPTER_URL/sync" \
  --oidc-service-account-email=csb-run-sa@rtp-transversal-dev.iam.gserviceaccount.com \
  --header="Content-Type: application/json" --body-file=payload.json

Rotar el token de HubSpot

printf '%s' "$TOKEN" | gcloud secrets versions add \
  csb-dev-hubspot-api-token --data-file=- --project=rtp-transversal-dev
# luego forzar nueva revisión (redeploy) para recargar la config

Desarrollo local y pruebas

Comandos

npm install
cp .env.example .env      # HUBSPOT_API_TOKEN, USE_SECRET_MANAGER=false
npm run dev               # nest start --watch → http://localhost:8080/docs
npm test                  # jest, *.spec.ts junto al código
npm run test:cov          # cobertura (umbral CI: 65 %)
npm run lint
npm run build             # dist/ incluye docs/public

Integración con mocks

Desde la carpeta padre rotoplas/, docker-compose.integration.yml levanta CSB + los tres adapters + mocks (SAP, HubSpot, Bringoz), Pub/Sub y Cloud Tasks emulados. El CRM adapter queda en http://localhost:8082 con Swagger en /docs/api. Ver INTEGRATION.md.

Pruebas existentes

Suites: transformer, hubspot-client (fetch mockeado, incluyendo 409 → search → PATCH), sync.service, sync.controller, result-publisher, http-exception.filter, health.controller, docs. La prueba "degraded when Secret Manager check fails" usa el cliente real de Secret Manager: en una máquina con credenciales de gcloud puede tardar y fallar por timeout; en Cloud Build pasa.

Convenciones

  • Ramas feature/* → dev → main. Nunca push directo a main.
  • Conventional Commits (feat:, fix:, chore:, infra:, docs:).
  • Secretos solo en Secret Manager o .env local ignorado por git.
  • Todo recurso GCP en Terraform.

Pendientes y deuda técnica

  • pendienteEndpoints de la saga order-to-hubspot: POST /contacts/resolve, POST /deals/transform, POST /deals/sync referenciados por el CSB y aún no implementados.
  • pendienteConvención de naming (revisión de Hugo, 26 ago 2026): rutas en inglés/REST bajo /crm/…, secretos con prefijo crm-<env>-… y repositorio de Artifact Registry propio crm-adapter. El código actual sigue usando csb- y csb-adapters.
  • deudaDos variables para el topic: RESULT_TOPIC se carga en la config pero ResultPublisherService lee PUBSUB_RESULT_TOPIC directamente. Unificar para evitar publicar en un topic equivocado.
  • deudaLogger estructurado no conectado: StructuredLoggerService existe pero main.ts no lo registra con app.useLogger(); los objetos pasan por el logger por defecto de NestJS.
  • deudaSin dead-letter: tras 5 intentos la tarea se descarta silenciosamente. Evaluar DLQ o publicación de evento de fallo.
  • fueraMarketing Hub, webhooks de HubSpot → CSB y Custom Objects son iniciativas separadas.
  • hechoMigración desde rtp-csb, Terraform propio, pipeline con canary (no-traffic + Ready check + promote), portal de documentación servido por el servicio.

Referencias

RecursoDónde
Repositoriogitlab.com/…/adapters/rtp-crm-adapter
Core Service Busgitlab.com/…/rtp-csb · web csb-dev.rotoplas.com
ERP Adapter (SAP)rtp-erp-adapter · Cloud Run csb-erp-adapter · csb-dev.rotoplas.com/adapters/erp/docs
Delivery Adapter (Bringoz/Nexus)rtp-delivery-adapter · Cloud Run csb-delivery-adapter-<env>-api · csb-dev.rotoplas.com/adapters/delivery/docs
Guía de integración y despliegue/docs/how-it-works — secrets, IAM, primer deploy, Terraform, smoke test E2E
API Reference/docs/api (Swagger UI) · /docs/openapi.json
Documentos del repoREADME.md, CLAUDE.md, SOUL.md, NEXT_STEPS.md, PIPELINE.md
HubSpotCRM API v3 — Contacts