How OpenCTI Connector Architecture Works with RabbitMQ Message Queues

OpenCTI uses RabbitMQ as a message bus where each connector maintains two dedicated queues—a "listen" queue for connector-to-platform communication and a "push" queue for platform-to-connector work distribution, managed through the rabbitmq.js module.

The OpenCTI connector architecture relies on RabbitMQ to decouple data ingestion, enrichment, and processing workflows. Each connector operates as an independent process that communicates with the platform exclusively through message queues, ensuring reliable asynchronous processing. This architecture is implemented in the OpenCTI-Platform/opencti repository, specifically within the GraphQL backend's RabbitMQ abstraction layer.

Core Architecture: Dual-Queue Design

Every connector in OpenCTI maintains two dedicated queues that isolate traffic by direction and purpose. This design prevents message collisions and enables independent scaling of producers and consumers.

The listen queue handles traffic from the connector to the platform:

  • Direction: Connector → Platform
  • Exchange: CONNECTOR_EXCHANGE (configured via amqp.connector.exchange)
  • Routing key: listen_routing_<connector-id>
  • Purpose: Publishes events such as new STIX objects, status updates, and processing confirmations

The push queue handles traffic from the platform to the connector:

  • Direction: Platform → Connector
  • Exchange: WORKER_EXCHANGE (configured via amqp.worker.exchange)
  • Routing key: push_routing_<connector-id>
  • Purpose: Distributes work items such as STIX bundles to ingest, enrichment requests, or remediation tasks

Queue Lifecycle and Registration

Queues are provisioned automatically when the platform initializes, ensuring connectors have immediate infrastructure available.

Automatic Queue Provisioning

The registerConnectorQueues function in opencti-platform/opencti-graphql/src/database/rabbitmq.js (lines 529-560) handles the complete setup:

  1. Asserts the CONNECTOR_EXCHANGE and WORKER_EXCHANGE exchanges
  2. Declares the listen and push queues for the connector
  3. Binds each queue to its respective exchange using the connector-specific routing keys
  4. Configures dead-letter routing for oversized bundles

Configuration Structure

The connectorConfig function (lines 16-28 in the same file) bundles all RabbitMQ connection parameters into a structured object. This includes queue names, routing keys, exchange identifiers, and dead-letter routing configuration used when messages exceed size limits.

Message Flow: Publishing and Consumption

Communication follows a strict request-response pattern mediated by the two exchanges.

Platform-to-Connector Delivery

When the platform needs to assign work to a connector, it calls pushToWorkerForConnector (lines 71-74 in rabbitmq.js). This function publishes to the WORKER_EXCHANGE using the connector's unique push routing key. The message typically contains a STIX bundle or an enrichment request.

Connector-to-Platform Responses

Connectors report results or stream data using pushToConnector (lines 76-78). This publishes to the CONNECTOR_EXCHANGE with the listen routing key, allowing the platform to consume status updates, newly created entities, or processing confirmations.

Message Consumption

Connectors initiate consumption by calling consumeQueue (lines 86-124). This function opens a channel on the listen queue, registers an asynchronous callback, and continuously processes platform-generated messages. The platform itself consumes push queue messages through internal worker managers such as taskManager and syncManager.

Internal Queue Types

Beyond external connectors, OpenCTI utilizes the same RabbitMQ infrastructure for internal background processing.

The platform generates dedicated queues for:

  • Background tasks – Created via getInternalBackgroundTaskQueues and limited by app:task_scheduler:max_queues_breakdown configuration. These handle asynchronous data processing jobs.
  • Playbooks – Generated by getInternalPlaybookQueues with one queue per stored playbook definition, created at runtime to orchestrate automation workflows.
  • Synchronization streams – Managed by getInternalSyncQueues to handle OpenCTI Streams for live data replication between instances.

All internal queues are registered using the same registerConnectorQueues logic, ensuring consistent durability, routing, and dead-letter handling across external and internal workers.

Monitoring and Health Checks

OpenCTI exposes several utilities to ensure the message bus remains healthy and to optimize workload distribution.

The rabbitMQIsAlive function (lines 60-66 in rabbitmq.js) performs periodic connectivity checks to verify the broker is reachable. For operational visibility, getConnectorQueueSize retrieves the current message count for a specific connector's push queue, enabling back-pressure detection.

The platform uses getBestBackgroundConnectorId to implement intelligent load balancing. By comparing queue depths across available background task connectors, the system routes new jobs to the worker with the lowest current backlog, preventing individual queue saturation.

Code Examples

The following snippets demonstrate the essential interactions for connector developers and platform integrators.

Building Connector Configuration

import { connectorConfig } from './database/rabbitmq';

const config = connectorConfig('my-connector-id');
console.log(config.listenQueue); // listen_routing_my-connector-id

Registering Queues Programmatically

await registerConnectorQueues(
  'my-connector-id',
  'My Connector',
  'external',               // type: external, internal, etc.
  ['observable', 'artifact'] // scopes this connector handles
);

Pushing Work to a Connector

import { pushToWorkerForConnector } from './database/rabbitmq';

await pushToWorkerForConnector('my-connector-id', { 
  action: 'ingest', 
  data: stixBundle 
});

Consuming Messages in a Connector

import { consumeQueue } from './database/rabbitmq';

await consumeQueue(context, 'my-connector-id', setConnection, async (ctx, msg) => {
  const payload = JSON.parse(msg);
  // Process platform-assigned work
  await processStixBundle(payload);
});

Responding to the Platform

import { pushToConnector } from './database/rabbitmq';

await pushToConnector('my-connector-id', { 
  status: 'processed', 
  entityId: createdEntity.id 
});

Key Source Files

Understanding the connector architecture requires familiarity with these specific files in the OpenCTI-Platform/opencti repository:

Summary

  • OpenCTI's connector architecture relies on RabbitMQ as the central message bus, ensuring reliable asynchronous communication between the platform and external connectors.
  • Each connector operates two dedicated queues: a listen queue (connector → platform) using CONNECTOR_EXCHANGE and a push queue (platform → connector) using WORKER_EXCHANGE.
  • Queue lifecycle management is handled by registerConnectorQueues in src/database/rabbitmq.js, which provisions exchanges, declares queues, and binds routing keys automatically.
  • Message flow uses pushToWorkerForConnector for platform-initiated work and pushToConnector for connector responses, with consumeQueue handling continuous message consumption.
  • Internal components like background tasks, playbooks, and sync streams reuse the same queue infrastructure, ensuring architectural consistency across the platform.

Frequently Asked Questions

How does OpenCTI ensure message delivery reliability between connectors and the platform?

OpenCTI implements dedicated exchanges and routing keys for each connector, with the registerConnectorQueues function configuring dead-letter routing for oversized bundles. The platform uses persistent message delivery and confirms queue declaration before accepting connector registrations, ensuring that messages are not lost during transmission.

What is the difference between the CONNECTOR_EXCHANGE and WORKER_EXCHANGE in OpenCTI?

The CONNECTOR_EXCHANGE (configured via amqp.connector.exchange) handles messages flowing from connectors to the platform via listen queues, while the WORKER_EXCHANGE (configured via amqp.worker.exchange) distributes work items from the platform to connectors via push queues. This separation prevents routing conflicts and enables independent scaling of ingestion and processing workloads.

Can internal OpenCTI components use the same queue infrastructure as external connectors?

Yes, internal components such as background tasks, playbooks, and synchronization streams utilize the same registerConnectorQueues logic found in src/database/rabbitmq.js. These internal queues are generated by functions like getInternalBackgroundTaskQueues and getInternalPlaybookQueues, ensuring consistent durability, routing, and monitoring across both external connectors and platform-native workers.

How does OpenCTI handle load balancing across multiple connector instances?

The platform uses the getBestBackgroundConnectorId function to monitor queue depths via getConnectorQueueSize and routes new jobs to the connector with the lowest current backlog. This intelligent load balancing prevents individual queue saturation and ensures optimal distribution of processing workloads across the connector pool.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →