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, and CONNECTOR_* variables from YAML or JSON configs.
  • Structured logging – Emitting JSON logs to stdout and the platform UI via self.log_info, self.log_debug, and self.log_error.
  • State persistence – Storing connector metadata (e.g., last run timestamps) in the OpenCTI KV store via get_state() and set_state().
  • Work lifecycle – Registering ingestion jobs with self.api.work.initiate_work and marking completion with to_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:

  1. Self-triggered (EXTERNAL_IMPORT, STREAM) – The connector runs on an interval using self.helper.listen(duration=self.interval, message_callback=self._process).
  2. 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 or listen_stream for SSE-based streams.
  • Persist state: Use get_state() and set_state() to store timestamps or offsets between runs, preventing duplicate processing.
  • Register work items: Always call helper.api.work.initiate_work before ingestion and to_processed after 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:

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 →