# STIX2 Bundle Processing Flow in PyCTI: End-to-End Technical Guide

> Explore the STIX2 bundle processing flow in PyCTI. Learn how OpenCTI connectors efficiently handle enrichment, validation, splitting, and dispatch for seamless data ingestion.

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

---

**The STIX2 bundle processing flow in PyCTI begins when `OpenCTIConnectorHelper.send_stix2_bundle` is called, proceeding through enrichment, validation, optional export, dependency-aware splitting, and finally dispatch to either a RabbitMQ queue or the OpenCTI API for ingestion.**

The Python client library `pycti` provides the primary interface for connectors to submit threat intelligence data to the OpenCTI platform. Understanding how it handles STIX2 bundles is essential for building reliable connectors and debugging ingestion issues. This guide traces the complete flow from the initial helper call through the final worker import, referencing the actual implementation in the `OpenCTI-Platform/opencti` repository.

## The 10-Step STIX2 Bundle Processing Flow

### 1. Helper Entry Point

The flow starts when a connector invokes `helper.send_stix2_bundle(bundle=<json>, **kwargs)` in [`opencti_connector_helper.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_connector_helper.py) at line 3390. This method serves as the central dispatcher that coordinates all subsequent processing stages.

### 2. Context Enrichment

For *enrichment* connectors, the helper injects shared organization references into every object of the bundle. This occurs in the `enrichment_shared_organizations` block between lines 63 and 98 of [`opencti_connector_helper.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_connector_helper.py), ensuring that enriched data carries the proper administrative context.

### 3. Playbook Shortcut

If the connector executes inside a *playbook step*, the bundle is sent directly to the playbook engine via `api.playbook.playbook_step_execution`, and the method returns early at lines 301-305. This bypasses the standard queue-based ingestion path.

### 4. Workbench and Draft Handling

When validation is required, the bundle is uploaded as a pending file (Workbench) or a draft is created at lines 310-328. The import aborts at this stage, deferring actual ingestion until a user validates the content through the OpenCTI interface.

### 5. Optional Directory and S3 Export

Before queue submission, the helper checks if the connector is configured to write bundles to a local directory or an S3 bucket. It builds a **message bundle** via `_create_message_bundle` and writes or exports it:

- **Directory export**: Lines 35-42 of [`opencti_connector_helper.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_connector_helper.py)
- **S3 export**: Lines 34-41 of [`opencti_connector_helper.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_connector_helper.py)

### 6. Bundle Splitting

The heavy-lifting of splitting large bundles into small, dependency-ordered chunks is performed by `OpenCTIStix2Splitter.split_bundle_with_expectations` in [`opencti_stix2_splitter.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_stix2_splitter.py) (lines 44-89). The splitter:

- Builds a **dependency graph** from `*_ref` and `*_refs` fields
- Removes incompatible items (e.g., relationships without `source_ref`/`target_ref`)
- Sorts objects by dependency count (least dependent first)
- Returns a list of ready-to-import bundles

### 7. Queue vs API Dispatch

The helper determines the transport method based on configuration:

- **Queue path (default)**: If `bundle_send_to_queue` is `True`, the helper opens a RabbitMQ (AMQP) connection and publishes each chunk as a **QUEUE_BUNDLE** message via `_send_bundle` (lines 361-398).
- **API path**: If the protocol is `api`, the entire original bundle is sent in a single request to `api.stix2.import_bundle` (lines 449-453).

### 8. RabbitMQ Message Format

Each chunk is wrapped in a JSON message containing:
- `bundle_type`
- `applicant_id`
- `action_sequence`
- `entities_types`
- Base-64 encoded STIX2 payload
- `update` flag
- Optional `draft_id` and `work_id`

This structure is defined at lines 700-711 of [`opencti_connector_helper.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_connector_helper.py).

### 9. Worker Consumption

The `opencti-worker` process receives the message via [`push_handler.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/push_handler.py) (lines 120-150). It validates the message, optionally re-splits extremely large bundles, and calls `api.stix2.import_bundle` to persist objects to the database.

### 10. Metrics and Error Handling

Throughout the flow, the helper updates OpenCTI metrics such as `bundle_send` and `error_count` (e.g., `self.metric.inc("bundle_send")` at line 726), logs progress, and implements retry logic for transient AMQP errors.

## Deep Dive: The STIX2 Splitter Algorithm

The `OpenCTIStix2Splitter` ensures that objects are imported in the correct order to satisfy STIX2 reference constraints. The algorithm executes as follows:

1. **Flatten the bundle**: Every object (including internal IDs from OpenCTI extensions) is stored in a flat `raw_data` dictionary keyed by `item["id"]`.
2. **Recursive enlist**: For each object, `enlist_element` walks through all `*_ref` and `*_refs` fields, building a reference graph (`cache_refs`) while avoiding cycles and optionally cleaning missing references.
3. **Compatibility check**: Objects lacking required fields (e.g., a relationship without `source_ref`/`target_ref`) are moved to `incompatible_items`.
4. **Dependency counting**: Each object receives a `nb_deps` count indicating how many other objects it depends on.
5. **Sorting**: Objects are sorted by `nb_deps` (least dependent first) to ensure referenced objects exist before dependent ones.
6. **Chunk creation**: Each object (or group of objects with the same dependency count) is wrapped into a minimal STIX2 bundle using `stix2_create_bundle`.

This implementation resides in [`opencti_stix2_splitter.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_stix2_splitter.py) at lines 44-89 and 443-471.

## Practical Code Examples

### Sending a Bundle from a Connector

```python
from pycti import OpenCTIConnectorHelper, OpenCTIStix2
import stix2

helper = OpenCTIConnectorHelper(config={})  # Auto-loads YAML/env config

stix2_api = OpenCTIStix2(helper.api)       # Low-level STIX helper

# Create a simple STIX2 Malware object and bundle it

malware = stix2.Malware(name="ExampleMalware", is_family=False)
bundle = stix2.Bundle(objects=[malware]).serialize()

# Send the bundle – triggers the full processing flow automatically

helper.send_stix2_bundle(bundle)

```

### Direct Bundle Splitting (Advanced)

```python
from pycti.utils.opencti_stix2_splitter import OpenCTIStix2Splitter

big_bundle_json = open("large_bundle.json").read()
splitter = OpenCTIStix2Splitter()
expectations, incompatible, chunks = splitter.split_bundle_with_expectations(
    bundle=big_bundle_json,
    use_json=True,
    cleanup_inconsistent_bundle=True,
)

print(f"Expectations (chunks): {expectations}")
print(f"Incompatible objects: {len(incompatible)}")

# Each item in `chunks` is a tiny STIX2 bundle ready for import

```

### Playbook Step Execution

```python

# Inside a connector running as a playbook step:

helper.playbook = {
    "event_id": "event-123",
    "execution_id": "exec-456",
    "step_id": "step-1",
    "playbook_id": "my-playbook",
}
helper.send_stix2_bundle(bundle)

# The helper calls `api.playbook.playbook_step_execution` and returns early

```

## Key Files and Implementation Locations

| File (relative to repository) | Purpose |
|-------------------------------|---------|
| [`client-python/pycti/connector/opencti_connector_helper.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/client-python/pycti/connector/opencti_connector_helper.py) | Public API for connectors; implements `send_stix2_bundle`, queue handling, validation, directory/S3 export, and metric logging. |
| [`client-python/pycti/utils/opencti_stix2_splitter.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/client-python/pycti/utils/opencti_stix2_splitter.py) | Dependency-aware bundle splitting logic; handles reference graphs and chunking. |
| [`opencti-worker/src/push_handler.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti-worker/src/push_handler.py) | Queue consumer that processes `QUEUE_BUNDLE` messages and triggers database import. |
| [`client-python/pycti/api/opencti_api_client.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/client-python/pycti/api/opencti_api_client.py) | Low-level API client for workbench uploads, draft creation, and direct bundle imports. |
| [`client-python/tests/cases/connectors.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/client-python/tests/cases/connectors.py) | Integration tests demonstrating typical `send_stix2_bundle` usage patterns. |
| [`client-python/tests/01-unit/utils/test_opencti_stix2_splitter.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/client-python/tests/01-unit/utils/test_opencti_stix2_splitter.py) | Unit tests verifying splitter expectations, dependency ordering, and edge case handling. |

## Summary

- The **STIX2 bundle processing flow** in PyCTI begins with `OpenCTIConnectorHelper.send_stix2_bundle` in [`opencti_connector_helper.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/opencti_connector_helper.py).
- **Enrichment connectors** can inject shared organization references before processing.
- **Playbook steps** bypass the standard queue and return immediately after calling the playbook API.
- **Validation modes** (Workbench/Draft) upload bundles as pending files rather than importing them directly.
- **Export options** allow writing bundles to local directories or S3 buckets before queue submission.
- **Bundle splitting** uses `OpenCTIStix2Splitter` to create dependency-ordered chunks, ensuring objects are imported before their references.
- **Transport** defaults to RabbitMQ (AMQP) via `_send_bundle`, with fallback to direct API calls.
- **Worker consumption** handles final validation, optional re-splitting, and database persistence via [`push_handler.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/push_handler.py).

## Frequently Asked Questions

### What is the entry point for sending STIX2 bundles in PyCTI?

The entry point is the `send_stix2_bundle` method of the `OpenCTIConnectorHelper` class, located at line 3390 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). This method orchestrates the entire flow from validation to final dispatch.

### How does PyCTI handle large STIX2 bundles that exceed message queue limits?

PyCTI handles large bundles through the `OpenCTIStix2Splitter` class in [`client-python/pycti/utils/opencti_stix2_splitter.py`](https://github.com/OpenCTI-Platform/opencti/blob/main/client-python/pycti/utils/opencti_stix2_splitter.py). The splitter constructs a dependency graph from object references, removes incompatible items, and divides the bundle into small, ordered chunks. Each chunk is processed independently, ensuring that referenced objects are always imported before dependent ones.

### What is the difference between queue-based and API-based bundle delivery?

Queue-based delivery (the default) splits the bundle into chunks and publishes them as `QUEUE_BUNDLE` messages to RabbitMQ via `_send_bundle` (lines 361-398), allowing the OpenCTI worker to process them asynchronously. API-based delivery sends the entire original bundle in a single synchronous request to `api.stix2.import_bundle` (lines 449-453), bypassing the message queue entirely.

### When should I use the Workbench or Draft validation modes?

Use Workbench or Draft modes when you require human validation before importing data. When these modes are enabled, `send_stix2_bundle` uploads the bundle as a pending file or creates a draft (lines 310-328) and aborts further processing. The platform stores the bundle for later review, and the actual import only occurs after a user manually approves the workbench entry or draft.