How to Build a Custom Python Connector Using the OpenCTIConnectorHelper Class
To build a custom Python connector for OpenCTI, instantiate the OpenCTIConnectorHelper class from the pycti SDK to handle configuration, logging, state management, and STIX bundle submission, then implement your processing logic and call the helper's listen method to start the event loop.
The OpenCTIConnectorHelper class is the core abstraction in the OpenCTI Python client that simplifies integration with the platform. Located in the OpenCTI-Platform/opencti repository, this helper encapsulates all communication primitives—authentication, message queuing, work tracking, and data serialization—allowing developers to focus on extraction and transformation logic rather than transport boilerplate.
Understanding the OpenCTIConnectorHelper Architecture
All connector capabilities are centralized in client-python/pycti/connector/opencti_connector_helper.py (lines 1875–1900). The helper acts as a long-running daemon that bridges external data sources or streams with the OpenCTI knowledge graph.
Key responsibilities include:
- Configuration management – Parsing
OPENCTI_URL,OPENCTI_TOKEN, andCONNECTOR_*variables from YAML or JSON configs. - Structured logging – Emitting JSON logs to stdout and the platform UI via
self.log_info,self.log_debug, andself.log_error. - State persistence – Storing connector metadata (e.g., last run timestamps) in the OpenCTI KV store via
get_state()andset_state(). - Work lifecycle – Registering ingestion jobs with
self.api.work.initiate_workand marking completion withto_processed. - Bundle submission – Chunking, validating, and transmitting STIX 2.1 bundles via
send_stix2_bundle. - Event consumption – Blocking loops for RabbitMQ (
listen) or Server-Sent Events (listen_stream) that dispatch messages to your callback.
Setting Up Connector Configuration
Connectors rely on environment variables or a config.yml file. The helper extracts values using the get_config_variable utility, which searches environment variables first, then falls back to nested config keys.
In client-python/pycti/connector/opencti_connector_helper.py (lines 1841–1845), the helper initializes core platform settings:
self.opencti_url = get_config_variable("OPENCTI_URL", ["opencti", "url"], config)
self.opencti_token = get_config_variable("OPENCTI_TOKEN", ["opencti", "token"], config)
self.connect_id = get_config_variable("CONNECTOR_ID", ["connector", "id"], config)
Your connector should load its own settings similarly:
from pycti import OpenCTIConnectorHelper, get_config_variable
import yaml
import os
config_path = os.path.join(os.path.dirname(__file__), "config.yml")
config = yaml.safe_load(open(config_path)) if os.path.isfile(config_path) else {}
helper = OpenCTIConnectorHelper(config)
api_url = get_config_variable("MY_API_URL", ["my_connector", "api_url"], config, True)
Implementing the Connector Class
Initialization and Helper Instantiation
The helper instance should be created once during connector startup. The class docstring (lines 1891–1897) demonstrates the canonical pattern:
helper = OpenCTIConnectorHelper(config)
This single object provides access to the OpenCTI API client (helper.api), pre-configured loggers, and scheduling metadata.
Processing Logic and Message Handling
Connectors operate in two primary modes:
- Self-triggered (
EXTERNAL_IMPORT,STREAM) – The connector runs on an interval usingself.helper.listen(duration=self.interval, message_callback=self._process). - On-demand (
INTERNAL_ENRICHMENT,INTERNAL_IMPORT) – The platform invokes the connector via RabbitMQ when specific events occur.
The listen method (lines 3062–3066) spawns a consumer thread that blocks until messages arrive, then delegates to your callback function:
def _process(self) -> str:
# Fetch and transform data
return "Message processed"
helper.listen(duration=300, message_callback=self._process)
State Management and Scheduling
For periodic connectors, persist state to avoid re-processing data. The helper provides KV store interactions and datetime utilities (lines 3028–3032):
import time
state = helper.get_state()
if state and "last_run" in state:
last_run = state["last_run"]
# Implement skip logic based on interval
helper.set_state({"last_run": int(time.time())})
next_run = helper.next_run_datetime(self.interval)
Submitting STIX Bundles to OpenCTI
After transforming external data into STIX 2.1 objects, submit them using send_stix2_bundle (lines 3390–3402). This method handles chunking large bundles, optional validation, and routing via the message queue or HTTP API.
Always register a work item before ingestion to enable progress tracking in the OpenCTI UI:
work_id = helper.api.work.initiate_work(
helper.connect_id,
f"Import run at {helper.current_timestamp()}"
)
helper.send_stix2_bundle(
bundle.serialize(),
entities_types=["Indicator", "Malware"],
update=True,
work_id=work_id
)
helper.api.work.to_processed(work_id, "Import completed successfully")
Complete External-Import Connector Example
Below is a minimal EXTERNAL_IMPORT connector that fetches indicators from a mock API and pushes them to OpenCTI every 60 seconds.
# src/main.py
import os
import yaml
import time
from pycti import OpenCTIConnectorHelper, get_config_variable
from stix2 import Indicator, Bundle
class MyIndicatorsConnector:
def __init__(self):
config_path = os.path.join(os.path.dirname(__file__), "config.yml")
config = yaml.safe_load(open(config_path)) if os.path.isfile(config_path) else {}
# Initialize the helper
self.helper = OpenCTIConnectorHelper(config)
# Connector-specific configuration
self.api_url = get_config_variable(
"MY_API_URL", ["my_connector", "api_url"], config, True
)
self.interval = int(
get_config_variable(
"MY_CONNECTOR_INTERVAL", ["my_connector", "interval"], config, True
)
)
def _fetch_indicators(self):
"""Simulate external API call."""
return [
{"value": "malicious.com", "type": "domain-name", "description": "C2 server"},
{"value": "192.168.1.1", "type": "ipv4-addr", "description": "Malicious host"},
]
def _process(self) -> str:
self.helper.log_info("Fetching indicators...")
raw = self._fetch_indicators()
stix_objects = [
Indicator(
name=item["value"],
description=item["description"],
indicator_types=["malicious-activity"],
pattern=f"[{item['type']}:value = '{item['value']}']",
valid_from="2024-01-01T00:00:00Z",
)
for item in raw
]
bundle = Bundle(objects=stix_objects).serialize()
work_id = self.helper.api.work.initiate_work(
self.helper.connect_id,
f"MyIndicators run @ {self.helper.current_timestamp()}"
)
self.helper.send_stix2_bundle(
bundle,
entities_types=["Indicator"],
update=True,
work_id=work_id,
)
self.helper.api.work.to_processed(work_id, "Bundle ingested")
self.helper.log_info("Import cycle complete")
return "Import completed"
def run(self) -> None:
# Self-triggered execution every self.interval seconds
self.helper.listen(duration=self.interval, message_callback=self._process)
if __name__ == "__main__":
connector = MyIndicatorsConnector()
connector.run()
Logging and Error Handling
The helper injects a pre-configured logger during initialization (lines 2267–2270). Use the convenience methods to ensure logs appear in both the container stdout and the OpenCTI platform:
self.helper.log_debug("Detailed debugging information")
self.helper.log_info("Standard operational message")
self.helper.log_error("Critical failure description")
These methods automatically include connector metadata and timestamps required for centralized observability.
Summary
- Instantiate once: Create a single
OpenCTIConnectorHelper(config)instance in your connector's__init__to manage platform authentication and configuration. - Configure with
get_config_variable: Retrieve environment-specific settings using the helper utility to ensure consistent fallback logic. - Choose the right listening mode: Use
listen(duration, callback)for self-triggered intervals orlisten_streamfor SSE-based streams. - Persist state: Use
get_state()andset_state()to store timestamps or offsets between runs, preventing duplicate processing. - Register work items: Always call
helper.api.work.initiate_workbefore ingestion andto_processedafter completion to track progress in the UI. - Submit via
send_stix2_bundle: Pass STIX 2.1 bundles to this method for automatic chunking, validation, and routing.
Frequently Asked Questions
What is the difference between listen and listen_stream in the OpenCTIConnectorHelper?
The listen method (lines 3062–3066) consumes messages from RabbitMQ and supports self-triggered intervals via the duration parameter, making it ideal for EXTERNAL_IMPORT connectors. The listen_stream method connects to the OpenCTI SSE (Server-Sent Events) endpoint for real-time streaming ingestion, which is used by STREAM connectors that react to live platform updates rather than polling external sources.
How does the helper handle configuration validation?
The get_config_variable function searches for values first in environment variables, then in the provided config dictionary using dot-notation paths (e.g., ["connector", "id"]). If a required variable is missing and no default is provided, the function raises an exception during initialization, ensuring the connector fails fast with a clear error message rather than running with invalid credentials.
Can I update existing OpenCTI objects when sending a bundle?
Yes. When calling send_stix2_bundle (lines 3390–3402), set the update=True parameter to enable upsert behavior. The helper will compare the incoming STIX objects against existing entities based on their STIX IDs and update attributes accordingly. You can also restrict which entity types are processed using the entities_types list to prevent unintended modifications to unrelated objects.
Where are the logging methods defined in the source code?
The shortcut logging methods (log_info, log_debug, log_error, etc.) are dynamically attached during the helper's __init__ in client-python/pycti/connector/opencti_connector_helper.py (lines 2267–2270). These methods wrap the underlying connector_logger instance to ensure all log entries include the connector ID and conform to the structured JSON format expected by the OpenCTI platform's logging aggregator.
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 →