OpenCTI Live Streaming: How to Consume Real-Time Threat Intelligence via SSE
OpenCTI's live streaming feature delivers real-time STIX 2.1 updates via Server-Sent Events (SSE), backed by a Redis stream that broadcasts create, update, delete, and merge events to authenticated clients.
The live streaming feature in OpenCTI enables security teams to consume Cyber Threat Intelligence (CTI) objects as they are created, modified, or deleted within the platform. As an open-source threat intelligence platform maintained by OpenCTI-Platform/opencti, it leverages a Redis stream and SSE middleware to push STIX 2.1 formatted data to connected clients in real-time. This architecture allows seamless integration with SIEMs, SOAR platforms, and custom analytics pipelines without polling overhead.
Architecture of the OpenCTI Live Streaming Feature
Core Components
The live streaming infrastructure relies on four primary components working in sequence:
-
Redis Stream (
stream.opencti): A persistent, ordered log defined insrc/database/stream/stream-utils.tsasLIVE_STREAM_NAME. All data layer operations publish JSON events to this stream viarawPushToStream. -
Stream Handler (
stream-handler.ts): Provides high-level helper functions includingstoreCreateEntityEvent,storeUpdateEvent,storeDeleteEvent, andstoreMergeEvent. These functions build STIX 2.1 payloads and callpushToStream(lines 31-41) to write events only whenpublishStreamEventflags permit and the operation is not in a draft context. -
SSE Middleware (
sseMiddleware.js): Exposes HTTP endpoints at/streamfor generic bypass streams and/stream/:idfor collection-specific feeds. TheliveStreamHandlerfunction authenticates requests, registers consumers, and initializes stream processors usingcreateStreamProcessor(lines 29-35 instream-handler.ts). -
Stream Collections: GraphQL objects managed in
src/resolvers/stream.jsandsrc/domain/stream.jsthat define filtered subsets of data usingStreamCollectionentities with custom filters and access controls.
Event Flow from Database to Client
When an object changes in OpenCTI, the platform propagates that change through the streaming stack:
-
Data Layer Change: A
store*function (e.g.,storeCreateEntity) triggersbuildCreateEventto generate a STIX 2.1 payload. -
Stream Publication: The handler calls
pushToStream, which writes the event to the Redis streamstream.openctivia the raw client implementation. -
Stream Processing: The
StreamProcessor(created viacreateStreamProcessor) continuously reads from Redis usingrawFetchStreamEventsRangeFromEventId. It resolves missing references if requested and validates events against filters usingisStixMatchFilterGroup. -
Client Delivery: Valid events pass through
client.sendEvent, which formats the SSE payload withid,event, anddatafields before transmission over the persistent HTTP connection.
Authentication and Access Control
Access to live streams is enforced through capability-based permissions and collection-specific restrictions:
-
Bypass Streams: Users with the
BYPASScapability can consume the generic stream at/api/streamwithout filter restrictions, receiving all platform events. -
Named Collections: Endpoints at
/api/stream/:idrequire authentication viaauthenticateForPublic. If a collection hasstream_public: true, any authenticated user may subscribe. Restricted collections require theTAXIIAPIcapability and membership in the collection'sauthorized_memberslist. -
Marking Restrictions: The middleware enforces data segregation by verifying that the requesting user possesses at least one of the markings specified in the collection's filters (see lines 51-63 in
sseMiddleware.js).
How to Consume OpenCTI Live Streams
Subscribing with cURL
Test connectivity to the generic live stream using standard HTTP clients:
curl -N -H "Authorization: Bearer <YOUR_TOKEN>" \
"http://localhost:4000/api/stream/live"
The -N flag disables buffering, ensuring each SSE data: line prints immediately upon arrival. Replace live with a specific Stream Collection UUID to receive filtered events.
Consuming from JavaScript
Use any SSE-compatible library to maintain a persistent connection:
// npm install eventsource
import EventSource from 'eventsource';
const token = '<YOUR_TOKEN>';
const collectionId = 'live'; // or a Stream Collection UUID
const url = `http://localhost:4000/api/stream/${collectionId}`;
const es = new EventSource(url, {
headers: { Authorization: `Bearer ${token}` },
});
es.onmessage = ({ data }) => {
const payload = JSON.parse(data);
console.log('Received STIX event:', payload);
};
es.onerror = (err) => {
console.error('Stream connection error:', err);
};
The payload contains STIX 2.1 objects generated by buildCreateEvent, buildUpdateEvent, or their deletion/merge counterparts.
Creating Filtered Stream Collections
Define tailored feeds through the GraphQL API to limit events by entity type, markings, or other criteria:
mutation CreateCollection($input: StreamCollectionAddInput!) {
streamCollectionAdd(input: $input) {
id
name
stream_public
filters
}
}
Variables:
{
"input": {
"name": "Public Indicators Only",
"description": "Real-time IOC feed",
"stream_public": true,
"filters": "{\"mode\":\"and\",\"filters\":[{\"key\":[\"entity_type\"],\"operator\":\"eq\",\"values\":[\"Indicator\"]}],\"filterGroups\":[]}"
}
}
After creation, clients connecting to /api/stream/<returned-id> receive only Indicator entities, with filtering occurring server-side in the StreamProcessor before SSE transmission.
Summary
- Redis-backed architecture: OpenCTI writes all CTI changes to a Redis stream (
stream.opencti) via functions instream-handler.ts, ensuring durable, ordered event delivery. - SSE protocol: Clients consume real-time updates through standard Server-Sent Events at
/api/streamendpoints, receiving STIX 2.1 formatted JSON. - Flexible filtering: Stream Collections allow administrators to expose specific data subsets using GraphQL-managed filters, while the processor handles resolution via
isStixMatchFilterGroup. - Strict access controls: The system enforces capabilities (
BYPASS,TAXIIAPI) and marking-based restrictions throughsseMiddleware.jsto prevent unauthorized data access. - Zero configuration: The platform automatically initializes the streaming stack via
initializeStreamStackwhen Redis is available, requiring no manual stream setup.
Frequently Asked Questions
What data format does the OpenCTI live streaming feature use?
The live streaming feature delivers events in STIX 2.1 format. When objects are created, updated, deleted, or merged, the platform generates standardized STIX bundles through functions like buildCreateEvent and buildUpdateEvent in stream-handler.ts, ensuring interoperability with other TIP and SOAR platforms.
How do I filter events by specific entity types or markings?
Create a Stream Collection via the GraphQL API (streamCollectionAdd mutation) and define JSON filters in the filters field. The StreamProcessor applies these filters using isStixMatchFilterGroup before transmitting events. For marking-based restrictions, ensure the collection specifies required markings and that users possess those markings in their profile.
What is the difference between the /stream and /stream/:id endpoints?
The /api/stream endpoint (often accessed as /api/stream/live) provides a generic bypass stream requiring the BYPASS capability and delivers all platform events without filtering. The /api/stream/:id endpoints connect to specific Stream Collections that apply custom filters and access controls, allowing granular public or restricted feeds tailored to specific use cases.
Do I need to manually configure Redis to enable live streaming?
No manual configuration is required. The platform automatically initializes the streaming infrastructure via initializeStreamStack in stream-handler.ts on startup, provided the REDIS_URL environment variable points to a running Redis instance. The stream stream.opencti is created automatically when the first event is published.
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 →