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 in src/database/stream/stream-utils.ts as LIVE_STREAM_NAME. All data layer operations publish JSON events to this stream via rawPushToStream.

  • Stream Handler (stream-handler.ts): Provides high-level helper functions including storeCreateEntityEvent, storeUpdateEvent, storeDeleteEvent, and storeMergeEvent. These functions build STIX 2.1 payloads and call pushToStream (lines 31-41) to write events only when publishStreamEvent flags permit and the operation is not in a draft context.

  • SSE Middleware (sseMiddleware.js): Exposes HTTP endpoints at /stream for generic bypass streams and /stream/:id for collection-specific feeds. The liveStreamHandler function authenticates requests, registers consumers, and initializes stream processors using createStreamProcessor (lines 29-35 in stream-handler.ts).

  • Stream Collections: GraphQL objects managed in src/resolvers/stream.js and src/domain/stream.js that define filtered subsets of data using StreamCollection entities 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:

  1. Data Layer Change: A store* function (e.g., storeCreateEntity) triggers buildCreateEvent to generate a STIX 2.1 payload.

  2. Stream Publication: The handler calls pushToStream, which writes the event to the Redis stream stream.opencti via the raw client implementation.

  3. Stream Processing: The StreamProcessor (created via createStreamProcessor) continuously reads from Redis using rawFetchStreamEventsRangeFromEventId. It resolves missing references if requested and validates events against filters using isStixMatchFilterGroup.

  4. Client Delivery: Valid events pass through client.sendEvent, which formats the SSE payload with id, event, and data fields 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 BYPASS capability can consume the generic stream at /api/stream without filter restrictions, receiving all platform events.

  • Named Collections: Endpoints at /api/stream/:id require authentication via authenticateForPublic. If a collection has stream_public: true, any authenticated user may subscribe. Restricted collections require the TAXIIAPI capability and membership in the collection's authorized_members list.

  • 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 in stream-handler.ts, ensuring durable, ordered event delivery.
  • SSE protocol: Clients consume real-time updates through standard Server-Sent Events at /api/stream endpoints, 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 through sseMiddleware.js to prevent unauthorized data access.
  • Zero configuration: The platform automatically initializes the streaming stack via initializeStreamStack when 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:

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 →