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 viaamqp.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 viaamqp.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:
- Asserts the
CONNECTOR_EXCHANGEandWORKER_EXCHANGEexchanges - Declares the listen and push queues for the connector
- Binds each queue to its respective exchange using the connector-specific routing keys
- 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
getInternalBackgroundTaskQueuesand limited byapp:task_scheduler:max_queues_breakdownconfiguration. These handle asynchronous data processing jobs. - Playbooks – Generated by
getInternalPlaybookQueueswith one queue per stored playbook definition, created at runtime to orchestrate automation workflows. - Synchronization streams – Managed by
getInternalSyncQueuesto 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:
-
opencti-platform/opencti-graphql/src/database/rabbitmq.js– Core RabbitMQ abstraction containingregisterConnectorQueues,connectorConfig,pushToWorkerForConnector,pushToConnector,consumeQueue, and health check utilities. -
opencti-platform/opencti-graphql/src/connector/connector-domain.ts– High-level domain logic exposing connector registration and queue discovery to other platform modules. -
opencti-platform/opencti-graphql/src/resolvers/connector.js– GraphQL resolver layer providing connector-related queries and mutations to the frontend API. -
opencti-platform/opencti-graphql/src/modules/playbook/playbook-domain.ts– Demonstrates internal queue usage for automation playbooks, showing howregisterConnectorQueuesapplies to internal components. -
opencti-platform/opencti-graphql/src/modules/pir/pir-domain.ts– Concrete implementation showing the Platform pushing jobs to a dedicated connector queue for PIR (Potential Incident Response) processing.
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_EXCHANGEand a push queue (platform → connector) usingWORKER_EXCHANGE. - Queue lifecycle management is handled by
registerConnectorQueuesinsrc/database/rabbitmq.js, which provisions exchanges, declares queues, and binds routing keys automatically. - Message flow uses
pushToWorkerForConnectorfor platform-initiated work andpushToConnectorfor connector responses, withconsumeQueuehandling 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →