How the OpenCTI Ingestion Manager Handles CSV Feeds: Architecture and Implementation
The OpenCTI ingestion manager processes CSV feeds by polling eligible ingestions every 30 seconds, downloading the source file, parsing rows through configurable mappers, converting them into STIX bundles, and queuing them for import while tracking state via SHA-256 hashes and execution timestamps.
The OpenCTI ingestion manager is a periodic background service that drives every built-in data feed type across the platform, including RSS, TAXII, JSON, and CSV sources. For CSV feeds specifically, the manager orchestrates a robust six-step pipeline that transforms raw delimited data into structured STIX objects through configurable mapping definitions. This article examines the implementation details in the OpenCTI-Platform/opencti repository, focusing on how src/manager/ingestionManager.ts and its supporting modules handle scheduling, throttling, parsing, and state persistence for CSV ingestion workflows.
Architecture Overview
The ingestion manager treats CSV feeds as scheduled background jobs that execute asynchronously. According to the source code in ingestionManager.ts, the pipeline follows six distinct phases:
- Selection: Query CSV ingestions where
ingestion_runningis true, checking scheduling periods and queue capacity. - Download: Fetch the raw CSV file using the configured URI and authentication via
fetchCsvFromUrl. - Parsing: Process the CSV with an inline or external mapper definition, handling headers and separators.
- Bundling: Transform each row into STIX objects using the CSV bundler and create work records.
- Queueing: Push generated bundles to the connector worker queue for asynchronous import.
- State Update: Record the file's SHA-256 hash, execution timestamp, and cursor position to enable resumable, idempotent processing.
Core Implementation in ingestionManager.ts
Scheduler and Entry Point
The manager initializes via initIngestionManager, which creates a setIntervalAsync timer that invokes ingestionHandler every SCHEDULE_TIME milliseconds (defaulting to 30 seconds). At lines 699–708, this scheduler ensures the service runs continuously as a background worker.
When the interval triggers, csvExecutor (starting at lines 77–84) builds a GraphQL filter to fetch all CSV ingestions where the ingestion_running flag is true:
const csvExecutor = async (context: AuthContext) => {
const filters = {
mode: 'and',
filters: [{ key: 'ingestion_running', values: [true] }],
filterGroups: [],
};
const opts = { filters, connectionFormat: false, noLimits: true };
const ingestions = await findAllCsvIngestion(context, SYSTEM_USER, opts);
// ...
};
Eligibility and Throttling
Before processing, the manager validates execution eligibility through two guards. First, isMustExecuteIteration (lines 72–78) checks whether the current time satisfies the ingestion's scheduling_period (e.g., "auto", "1h", "24h"). Second, shouldExecuteIngestion enforces the CSV_FEED_MIN_INTERVAL_MINUTES constraint to prevent overwhelming external sources.
The code also verifies queue pressure via queueDetails, ensuring messages_number is zero before fetching new data:
for (const ingestion of ingestions) {
if (isMustExecuteIteration(ingestion)) {
const { messages_number } = await queueDetails(connectorIdFromIngestId(ingestion.id));
if (messages_number === 0 && shouldExecuteIngestion(ingestion, CSV_FEED_MIN_INTERVAL_MINUTES)) {
await csvDataHandler(context, ingestion);
}
}
}
Data Retrieval and Parsing
The csvDataHandler function (lines 63–68) resolves the CSV mapper configuration—either an inline JSON definition or a reference to an external csvMapper entity—then downloads the file via fetchCsvFromUrl from src/utils/http-client.ts.
Subsequently, processCsvLines (lines 19–53) handles the raw CSV stream. It strips headers if configured, initializes a work record for tracking, and constructs a CsvBundlerIngestionOpts object containing the mapper, applicant user, and connector ID.
Bundling and Queue Distribution
The core transformation occurs in generateAndSendBundleProcess, implemented in src/parser/csv-bundler.ts. The bundler iterates over CSV rows and converts them into STIX objects (typically Indicators) based on the mapper's column-to-field mappings:
const { bundleCount, objectCount } = await generateAndSendBundleProcess(
ctx,
csvLines,
{
workId,
applicantUser,
csvMapper,
connectorId,
entity: undefined,
}
);
After bundling, the manager pushes the resulting STIX bundle to the connector worker queue:
await pushBundleToConnectorQueue(context, ingestion, {
type: 'bundle',
spec_version: '2.1',
id: `bundle--${uuidv4()}`,
objects,
});
State Management and Error Handling
Upon successful completion, processCsvLines updates the ingestion record (lines 55–58) with:
- The SHA-256 hash of the downloaded file to detect changes
- The
added_after_startcursor for incremental processing - The
last_execution_datetimestamp
Error handling occurs in the csvExecutor catch block (lines 4–7). Any failure during download or processing triggers patchCsvIngestion to update the last_execution_date, ensuring the minimum interval is respected before the next retry, while logging the failure for observability.
Supporting Modules and File Structure
The ingestion manager delegates specialized tasks to dedicated modules:
src/modules/ingestion/ingestion-csv-domain.ts: Provides CRUD helpers includingfindAllCsvIngestionandpatchCsvIngestionfor persistence operations.src/parser/csv-bundler.ts: ContainsgenerateAndSendBundleProcess, the core logic for transforming CSV rows into STIX 2.1 bundles.src/modules/internal/csvMapper/csvMapper-utils.ts: Parses mapper definitions, handling header detection, separator configuration, and column mapping.src/modules/internal/csvMapper/csvMapper-domain.ts: Manages external CSV mapper entities stored in the database.src/utils/http-client.ts: ImplementsfetchCsvFromUrlwith proper timeout, header management, and authentication handling.
Practical Configuration Example
To create a CSV ingestion via the GraphQL API:
mutation CreateCsvIngestion($input: IngestionCsvAddInput!) {
ingestionCsvAdd(input: $input) {
id
name
uri
csv_mapper_type
csv_mapper
}
}
With variables:
{
"input": {
"name": "Blocklist Feed",
"uri": "https://lists.blocklist.de/lists/all.txt",
"authentication_type": "none",
"csv_mapper_type": "inline",
"csv_mapper": {
"has_header": false,
"separator": ",",
"skip_line_char": "",
"attributes": [
{ "column_name": "value", "stix_field": "pattern", "default_value": "value: {value}" }
]
},
"scheduling_period": "auto",
"ingestion_running": true
}
}
Once created, the ingestion appears in the manager's next 30-second cycle, or can be triggered immediately via the /ingestion/manager/start endpoint.
Summary
- The OpenCTI ingestion manager polls for active CSV feeds every 30 seconds using
setIntervalAsynciningestionManager.ts. - Eligibility checks include scheduling periods, queue capacity (
messages_number), and minimum intervals (CSV_FEED_MIN_INTERVAL_MINUTES). - Data retrieval uses
fetchCsvFromUrlwith support for authentication headers and configurable mappers. - Transformation occurs via
generateAndSendBundleProcessin the CSV bundler, converting rows to STIX objects. - State tracking relies on SHA-256 hashes and timestamps to prevent duplicate processing and enable failure recovery.
- Error handling updates execution dates to maintain throttling constraints while logging failures.
Frequently Asked Questions
How often does the OpenCTI ingestion manager check for new CSV data?
The manager creates a setIntervalAsync timer in initIngestionManager that executes every 30 seconds (SCHEDULE_TIME). However, individual CSV ingestions only run if they satisfy their scheduling_period and respect the CSV_FEED_MIN_INTERVAL_MINUTES threshold since their last execution.
What prevents the same CSV file from being processed multiple times?
The manager calculates a SHA-256 hash of each downloaded file and stores it in the ingestion record. Before processing, it compares the current file hash against the stored value. Additionally, the added_after_start cursor tracks incremental feeds, ensuring only new rows are processed during subsequent runs.
Can OpenCTI handle CSV feeds without column headers?
Yes. The mapper configuration includes a has_header boolean flag. When set to false, the parser treats every line as data rows, mapping columns by index according to the mapper's column_name definitions. This is configured via csvMapper-utils.ts during the parsing phase.
Where are CSV ingestion errors logged in OpenCTI?
Errors occurring within csvExecutor or csvDataHandler are caught in the executor's catch block, which invokes patchCsvIngestion to update the last_execution_date and prevent immediate retries. Detailed error messages are logged through the standard OpenCTI logging infrastructure, accessible via the platform's logging configuration or the connector status dashboard updated by updateBuiltInConnectorInfo.
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 →