Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
{
"event": "create",
"id": "609198385",
"retry": null,
"data": {
"data": {
"id": "indicator--32c6886b-68e6-45e5-bb0f-0f5defa8b7ef",
"spec_version": "2.1",
"type": "indicator",
"extensions": {
"extension-definition--ea279b3e-5c71-4632-ac08-831c66a786ba": {
"extension_type": "property-extension",
"id": "15926f6f-44b9-446d-94e0-cf66c6aceec7",
"type": "Indicator",
"created_at": "2027-03-11T14:22:05.000Z",
"updated_at": "2027-03-11T14:39:05.000Z",
"is_inferred": false,
"creator_ids": [
"7d9118e3-5f2f-453d-a20a-b4acd9a50e7a"
],
"detection": false,
"score": 50,
"main_observable_type": "StixFile",
"observable_values": [
{
"type": "StixFile",
"hashes": {
"MD5": "c2f7c2f12563246ea8ce5a4ea425de8e"
}
},
{
"type": "StixFile",
"hashes": {
"SHA-256": "b59444a310faa5ab8a774161cd8961012e3c02daa70273bc75de8acfd8f504b7"
}
}
]
},
"extension-definition--322b8f77-262a-4cb8-a915-1e441e00329b": {
"extension_type": "property-extension"
}
},
"created": "2027-03-11T14:22:05.000Z",
"modified": "2027-03-11T14:39:05.000Z",
"revoked": false,
"confidence": 100,
"lang": "en",
"name": "sample-indicator",
"pattern": "[file:hashes.MD5 = 'c2f7c2f12563246ea8ce5a4ea425de8e' AND file:hashes.'SHA-256' = 'b59444a310faa5ab8a774161cd8961012e3c02daa70273bc75de8acfd8f504b7']",
"pattern_type": "stix",
"valid_from": "2027-03-11T14:20:05.000Z",
"valid_until": "2027-12-26T14:22:05.000Z"
},
"message": "creates a Indicator `sample-indicator`",
"origin": {
"referer": "init-create"
},
"version": "4"
}
}
159 changes: 154 additions & 5 deletions stream/pan-cortex-xdr-intel/src/connector/connector.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,49 @@
from __future__ import annotations

from typing import TYPE_CHECKING
import json
from typing import TYPE_CHECKING, Any, Protocol

from connector.models import EventIndicator
from pydantic import ValidationError

if TYPE_CHECKING:
from connector.settings import ConnectorSettings
from cortex_xdr_client import CortexXdrClient
from pycti import OpenCTIConnectorHelper


_SUPPORTED_EVENTS = {"create", "update", "delete"}
_SUPPORTED_ENTITY_TYPES = {"indicator"}
_SUPPORTED_OBSERVABLE_TYPES = {
"domain-name",
"hostname",
"ipv4-addr",
"ipv6-addr",
"stixfile",
}

_OPENCTI_OBSERVABLE_TYPES_TO_XDR_IOC_TYPES = {
"domain-name": "DOMAIN_NAME",
"hostname": "DOMAIN_NAME",
"ipv4-addr": "IP",
"ipv6-addr": "IP",
"stixfile": {"name": "FILENAME", "hash": "HASH"},
# XDR IOC types `PATH` and `MIXED` are not mapped for now
}


class StreamMessage(Protocol):
"""Type for the SSE message passed to `listen_stream` callbacks.

Only the attributes actually consumed by this connector are declared,
decoupling us from `filigran_sseclient.sseclient.Event`'s concrete shape
(and from adding it as an explicit dependency just for typing).
"""

event: str
data: str


class Connector:
"""
Cortex XDR Intel stream connector.
Expand Down Expand Up @@ -38,11 +74,124 @@ def __init__(
self.settings = settings
self.client = client

def _process_message(self, message: dict) -> None:
# TODO: implement IOC upsert/delete mapping from the stream message (#7184/#7185)
self.helper.connector_logger.info(
"Received stream event (baseline wiring phase)"
# Reusable exit message for fatal errors logging
self._exit_message = (
"Connector will exit to avoid further errors and/or exhausting the stream.\n"
"Please check connector's logs and report/fix the issue before restarting."
)

def _parse_indicator(self, data: dict[str, Any]) -> EventIndicator:
"""Build a minimal, internal representation of an Indicator's stream
`data` payload. Field casting/validation is delegated to `EventIndicator`.
"""
# Use `observable_values` extension provided by OpenCTI to extract the list of observables
observable_values = (
self.helper.get_attribute_in_extension("observable_values", data) or []
)
# Filter out observable types not supported by Cortex XDR Client
supported_observable_values = [
observable_value
for observable_value in observable_values
if observable_value.get("type", "").lower() in _SUPPORTED_OBSERVABLE_TYPES
]

# Parse and validate the indicator in stream `data` payload
return EventIndicator(
id=self.helper.get_attribute_in_extension("id", data),
description=data.get("description"),
observables=supported_observable_values,
valid_until=data.get("valid_until"),
score=self.helper.get_attribute_in_extension("score", data),
)

def _process_message(self, msg: StreamMessage) -> None:
"""Process a single stream event message.

Unsupported event or entity types are logged as a warning and
skipped. Any failure while decoding, parsing, or validating the
message (JSON decode error, `Indicator` validation error, or any
other unexpected exception) is logged with context and re-raised,
deliberately letting `pycti` kill the connector process rather than
risk silently missing or corrupting further events.
"""
event = msg.event
if event not in _SUPPORTED_EVENTS:
self.helper.connector_logger.warning(
"Unsupported event type, skipping it",
{"event": event},
)
return

try:
message_data = json.loads(msg.data)
except json.JSONDecodeError as err:
# This should never happen and if it does, it indicates a breaking change in `pycti`.
# To avoid data loss, the connector must stop and the issue must be investigated and fixed before resuming.
self.helper.connector_logger.error(
f"Failed to parse stream event's `data` payload as JSON.\n{self._exit_message}",
{
"event": event,
"data": msg.data,
"error": err,
},
)
raise # let `pycti` kill the connector process

entity_data = message_data.get("data", {})
entity_type = entity_data.get("type")
if entity_type not in _SUPPORTED_ENTITY_TYPES:
self.helper.connector_logger.warning(
"Unsupported entity type, skipping it",
{
"event": event,
"entity_type": entity_type,
},
)
return

try:
indicator = self._parse_indicator(entity_data)

self.helper.connector_logger.info(
"Parsed observable(s) from stream event",
{
"event": event,
"indicator_id": indicator.id,
"observables_count": len(indicator.observables),
},
)
except ValidationError as err:
# This should never happen and if it does, it indicates a breaking change in `pycti`.
# To avoid data loss, the connector must stop and the issue must be investigated and fixed before resuming.
self.helper.connector_logger.error(
f"Failed to parse indicator and/or observables from stream event.\n{self._exit_message}",
{
"event": event,
"entity_data": entity_data,
"error": err,
},
)
raise # let `pycti` kill the connector process

try:
if event in {"create", "update"}:
pass # TODO: upsert in Cortex XDR
elif event == "delete":
pass # TODO: delete from Cortex XDR

except Exception as err:
# Repetitive unexpected errors could consume the stream in vain (no action performed).
# To avoid data loss, the connector must stop and the issue must be investigated and fixed before resuming.
self.helper.connector_logger.error(
f"Unexpected error while processing stream event.\n{self._exit_message}",
{
"event": event,
"entity_data": entity_data,
"error": err,
},
)
raise # let `pycti` kill the connector process

def start(self) -> None:
"""Start the connector's main loop: listen to the OpenCTI stream and process each message."""
self.helper.listen_stream(self._process_message)
62 changes: 62 additions & 0 deletions stream/pan-cortex-xdr-intel/src/connector/models.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
from __future__ import annotations

from datetime import datetime
from typing import Any

from pydantic import BaseModel, ConfigDict, Field, field_validator


class IndicatorObservable(BaseModel):
"""A single observable extracted from an Indicator's `observable_values`."""

model_config = ConfigDict(frozen=True)

type: str
value: str


class EventIndicator(BaseModel):
"""Minimal, internal representation of an Indicator's stream `data` payload.

This is a plain data carrier meant to decouple downstream handlers
(upsert/delete) from the raw OpenCTI STIX payload shape. Only `data` is
represented here: the stream action (create/update/delete) is a separate,
stream-level concern and is passed alongside this event, not stored on it.
Fields are cast/validated by pydantic to fail fast and minimize downstream
development/runtime errors.
"""

model_config = ConfigDict(frozen=True)

id: str
description: str | None = Field(default=None)
observables: list[IndicatorObservable] = Field(default_factory=list)
valid_until: datetime | None = Field(default=None)
score: int | None = Field(default=None)

@field_validator("observables", mode="before")
def _validate_observables(
cls, value: list[dict[str, Any]] | None
) -> list[dict[str, Any]]:
"""Extract `observables` from the raw `observable_values` list returned by
`pycti.OpenCTIConnectorHelper.get_attribute_in_extension("observable_values", indicator)`.

For `StixFile` observables, the filename (`name`) and each hash
algorithm value are extracted as separate observables, since Cortex
XDR treats them as distinct IOC types (`FILENAME` vs `HASH`).
"""
if not value:
return []

observables = []
for observable_data in value:
observable_type = str(observable_data.get("type", ""))
if observable_type.lower() == "stixfile":
if name := observable_data.get("name"):
observables.append({"type": observable_type, "value": str(name)})
for value in (observable_data.get("hashes") or {}).values():
observables.append({"type": observable_type, "value": str(value)})
elif value := observable_data.get("value"):
observables.append({"type": observable_type, "value": str(value)})

return observables
Loading
Loading