# How to Build a Custom Python Connector Using the OpenCTIConnectorHelper Class

> Easily build a custom Python connector for OpenCTI using the OpenCTIConnectorHelper class. Streamline configuration, logging, and STIX bundle submission for seamless integration.

- Repository: [OpenCTI Platform/opencti](https://github.com/opencti-platform/opencti)
- Tags: how-to-guide
- Published: 2026-02-19

---

**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`](https://github.com/OpenCTI-Platform/opencti/blob/main/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`](https://github.com/OpenCTI-Platform/opencti/blob/main/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`](https://github.com/OpenCTI-Platform/opencti/blob/main/client-python/pycti/connector/opencti_connector_helper.py) (lines 1841–1845), the helper initializes core platform settings:

```python
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:

```python
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:

```python
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:

```python
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):

```python
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:

```python
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.

```python

# 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:

```python
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`](https://github.com/OpenCTI-Platform/opencti/blob/main/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.