STIX2 Bundle Processing Flow in PyCTI: End-to-End Technical Guide
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 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, 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 - S3 export: Lines 34-41 of
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 (lines 44-89). The splitter:
- Builds a dependency graph from
*_refand*_refsfields - 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_queueisTrue, 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 toapi.stix2.import_bundle(lines 449-453).
8. RabbitMQ Message Format
Each chunk is wrapped in a JSON message containing:
bundle_typeapplicant_idaction_sequenceentities_types- Base-64 encoded STIX2 payload
updateflag- Optional
draft_idandwork_id
This structure is defined at lines 700-711 of opencti_connector_helper.py.
9. Worker Consumption
The opencti-worker process receives the message via 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:
- Flatten the bundle: Every object (including internal IDs from OpenCTI extensions) is stored in a flat
raw_datadictionary keyed byitem["id"]. - Recursive enlist: For each object,
enlist_elementwalks through all*_refand*_refsfields, building a reference graph (cache_refs) while avoiding cycles and optionally cleaning missing references. - Compatibility check: Objects lacking required fields (e.g., a relationship without
source_ref/target_ref) are moved toincompatible_items. - Dependency counting: Each object receives a
nb_depscount indicating how many other objects it depends on. - Sorting: Objects are sorted by
nb_deps(least dependent first) to ensure referenced objects exist before dependent ones. - 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 at lines 44-89 and 443-471.
Practical Code Examples
Sending a Bundle from a Connector
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)
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
# 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 |
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 |
Dependency-aware bundle splitting logic; handles reference graphs and chunking. |
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 |
Low-level API client for workbench uploads, draft creation, and direct bundle imports. |
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 |
Unit tests verifying splitter expectations, dependency ordering, and edge case handling. |
Summary
- The STIX2 bundle processing flow in PyCTI begins with
OpenCTIConnectorHelper.send_stix2_bundleinopencti_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
OpenCTIStix2Splitterto 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.
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. 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. 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.
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 →