How the OpenCTI Sync Manager Handles Data Synchronization: Architecture and Event Processing

The OpenCTI Sync Manager continuously consumes remote Server-Sent Events (SSE) streams, transforms each incoming event into internal create, update, or delete actions, and ensures fault-tolerant synchronization across clustered instances through Redis-backed locking and resumable state persistence.

The synchronization capability in OpenCTI enables real-time data sharing between platform instances by consuming remote streams and applying changes locally. This article examines how the Sync Manager in the OpenCTI-Platform/opencti repository orchestrates this process, from cluster-wide coordination to per-stream event handling.

Core Architecture of the OpenCTI Sync Manager

The manager operates as a singleton service within a cluster, utilizing a layered architecture that separates scheduling, coordination, and stream processing concerns.

Cluster Coordination via Redis Locking

To prevent duplicate processing across horizontally scaled API nodes, the manager acquires a global Redis lock before activation. In opencti-platform/opencti-graphql/src/manager/syncManager.js, the constant SYNC_MANAGER_KEY (defaulting to sync_manager_lock) defines the lock resource. The lockResources([SYNC_MANAGER_KEY]) call ensures only one instance per cluster executes the processing loop at any time. If another node holds the lock, the handler exits silently without error.

The Processing Loop and Scheduler

Once the lock is acquired, initSyncManager creates an asynchronous scheduler that executes syncManagerHandler every SCHEDULE_TIME milliseconds (default 10,000ms). Each iteration fetches all Sync entities from the database using topEntitiesList, then manages their lifecycle based on the running flag. The system starts new syncs, stops paused ones, and cleans up deleted configurations by stopping their associated manager instances.

Per-Sync Instance Management

For each active sync configuration, the system instantiates a dedicated handler via syncManagerInstance(id). This factory function returns an object with start() and stop() methods that manage a single remote stream connection. The per-sync instance maintains its own HTTP-SSE connection, event parsing logic, and retry mechanisms, isolating stream-specific failures from the global scheduling loop.

Event Processing Flow

When a per-sync instance starts, it enters a reconnection loop that handles the full lifecycle of remote events.

Stream Connection and Event Types

The start() method builds an SSE URI using createSyncHttpUri with the last known state (lastState or lastEventDate), then calls streamEvents(sseUri, sync)—an async generator that feeds the raw HTTP stream through eventsource-parser. The system handles four distinct event types:

  • connected: Stores the connection identifier for tracking.
  • heartbeat: Updates state timestamp without processing payload data.
  • consumer_metrics: Persists performance metrics to the database.
  • data: Triggers the full transformation and ingestion pipeline.

Data Transformation and Worker Dispatch

For data events, the manager calls transformDataWithReverseIdAndFilesData to download attached files and apply reverse-patch operations on STIX identifiers. The enriched payload is base-64 encoded and dispatched to the internal worker queue via pushToWorkerForConnector. This decouples stream consumption from entity creation, allowing the worker pool to handle the actual database writes asynchronously.

State Persistence and Resume Capability

After successfully processing any event (including heartbeats), the manager invokes saveCurrentState to persist current_state_date to the Sync entity in the database. This timestamp enables resumable connections—if the stream disconnects or the platform restarts, the manager reconnects using the last processed timestamp, ensuring no data loss during transient failures.

Configuration and Enablement

The manager respects three primary configuration keys defined in opencti-platform/opencti-graphql/src/config/conf.ts:

  • sync_manager:enabled: Boolean toggle (default false) controlling whether the feature is active.
  • sync_manager:interval: Scheduler frequency in milliseconds (default 10000).
  • sync_manager:lock_key: Redis key for cluster coordination (default sync_manager_lock).

Administrators can override these via environment variables or the platform settings UI.

Integration with the OpenCTI Frontend

The frontend dynamically exposes synchronization controls based on backend capability. In opencti-platform/opencti-front/src/private/components/data/Sync.tsx, the application checks platformModuleHelpers.isSyncManagerEnable() to conditionally render the Streams UI. This ensures administrators only see sync configuration options when the backend manager is enabled and ready to process streams.

Practical Code Examples

Starting and Monitoring the Manager

While the platform bootstraps the manager automatically, manual interaction is possible:

import syncManager from './src/manager/syncManager';

// Start the scheduler and acquire cluster lock
await syncManager.start();

// Check operational status
const status = syncManager.status();
console.log(status); // { id: 'SYNC_MANAGER', enable: true, running: true }

Graceful Shutdown

During platform termination, the manager supports clean shutdown:

await syncManager.shutdown();
// Clears the scheduler interval and stops all active per-sync instances

Manual Per-Sync Instance Creation

For testing or custom implementations, instantiate individual sync handlers:

import { syncManagerInstance } from './src/manager/syncManager';

const sync = syncManagerInstance('sync-uuid-here');
await sync.start(executionContext('test'));
// ... processing occurs ...
await sync.stop();

Summary

  • The OpenCTI Sync Manager uses Redis distributed locking (SYNC_MANAGER_KEY) to ensure singleton execution across clustered nodes.
  • A scheduler loop (processStep) runs every 10 seconds to manage the lifecycle of configured syncs.
  • Per-sync instances handle HTTP-SSE connections, parsing events via streamEvents and dispatching work through pushToWorkerForConnector.
  • State persistence via saveCurrentState enables resumable streams after disconnections or restarts.
  • The feature is disabled by default and must be enabled via sync_manager:enabled configuration.

Frequently Asked Questions

How does the Sync Manager prevent duplicate processing in clustered deployments?

The manager acquires a Redis lock using lockResources([SYNC_MANAGER_KEY]) before executing the processing loop. Only the node holding the lock runs the scheduler; others exit silently. This guarantees a single active coordinator per cluster regardless of API node count.

What happens if the connection to a remote stream is interrupted?

The per-sync instance maintains a reconnection loop (while (running)) that automatically retries connections every 5 seconds after failures. Because saveCurrentState persists the last processed timestamp after each event, the manager reconnects using the lastEventDate parameter, resuming exactly where it left off without reprocessing historical data.

How are file attachments handled during synchronization?

When processing data events, the manager invokes transformDataWithReverseIdAndFilesData to download file attachments from the remote instance and apply identifier reverse-patching. This ensures that file-based observables and reports are correctly transferred and referenced in the local database before the worker processes the entity creation.

Can the Sync Manager be enabled or disabled at runtime?

While the sync_manager:enabled configuration is read at startup, the manager exports a status() method that reflects the current operational state. Frontend components check platformModuleHelpers.isSyncManagerEnable() to conditionally render the UI, ensuring that administrators cannot configure streams when the backend capability is disabled. Runtime toggling requires a platform restart to reinitialize the scheduler.

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 →