Skip to content
Closed
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
Expand Up @@ -101,7 +101,7 @@ state.attach_opencti_connector_helper(helper) # establish the connection with Op
state.load()

if state.last_run:
self.helper.connector_logger.info("Last run:", {"last_run": state.last_run})
logger.info("Last run:", {"last_run": state.last_run})

state.last_run = datetime.now(tz=timezone.utc)
state.save()
Expand Down Expand Up @@ -203,3 +203,8 @@ The `states` module has no dependency on other SDK modules (`settings`, `models`
- Related TDR: [Typing and validation of configurations with Pydantic Settings](https://github.com/OpenCTI-Platform/connectors/blob/master/connectors-sdk/TDRs/2025-10-01-Typing_and_validation_of_configurations_with_Pydantic_Settings.md)
- [Pydantic BaseModel documentation](https://docs.pydantic.dev/latest/concepts/models/)
- [Pydantic serialization (JSON mode)](https://docs.pydantic.dev/latest/concepts/serialization/#json-mode)


## Update (2026-05-28)

Update logging in code examples. No technical changes.
13 changes: 6 additions & 7 deletions connectors-sdk/connectors_sdk/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,13 @@

__version__ = "0.1.0"

from connectors_sdk.connectors.external_import._work_manager import WorkManager
from connectors_sdk.connectors.external_import.base_data_processor import (
BaseDataProcessor,
)
from connectors_sdk.connectors.external_import.external_import_connector import (
ExternalImportConnector,
)
from connectors_sdk.connectors.external_import.logger import ConnectorLogger
from connectors_sdk.logging.logger import Logger
from connectors_sdk.settings.annotated_types import (
DatetimeFromIsoString,
ListFromString,
Expand All @@ -34,6 +33,11 @@
from connectors_sdk.states.states import ExternalImportConnectorState

__all__ = [
# Base Connectors
"BaseDataProcessor",
"ExternalImportConnector",
# Logger
"Logger", # mostly for typing purposes
# Base Settings
"BaseConnectorSettings",
# Base Configs
Expand All @@ -54,9 +58,4 @@
"DeprecatedField",
# Connector States
"ExternalImportConnectorState",
# Connector base classes
"ExternalImportConnector",
"ConnectorLogger",
"BaseDataProcessor",
"WorkManager",
]
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
"""Module containing base classes for connector development.

This module provides the foundational classes for building OpenCTI connectors:
- ConnectorLogger: Logging wrapper to avoid direct pycti dependency
- Logger: Logging wrapper to avoid direct pycti dependency
- WorkManager: Work lifecycle management (initiate, send bundles, complete)
- BaseDataProcessor: Abstract base class for data collection/processing
- ExternalImportConnector: Full orchestration for external import connectors
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,14 @@

from __future__ import annotations

from typing import Any
from typing import TYPE_CHECKING, Any, ClassVar

from connectors_sdk.connectors.external_import.logger import ConnectorLogger
from connectors_sdk.logging.sdk_logger import sdk_logger
from pycti import OpenCTIConnectorHelper

if TYPE_CHECKING:
from connectors_sdk.logging._base_logger import BaseLogger


class _Work:
"""Manage the lifecycle of a single work in OpenCTI.
Expand All @@ -37,54 +40,49 @@ class _Work:
name: The human-readable name of the work, displayed in the OpenCTI UI.
"""

logger: ClassVar[BaseLogger] = sdk_logger.get_child("WorkManager._Work")

def __init__(
self,
helper: OpenCTIConnectorHelper,
work_id: str,
work_name: str,
helper: OpenCTIConnectorHelper,
logger: ConnectorLogger,
) -> None:
"""Initialize the work context.

Args:
helper: The ``OpenCTIConnectorHelper`` instance.
work_id: The work ID returned by OpenCTI.
work_name: The human-readable name of the work, displayed in the OpenCTI UI.
helper: The ``OpenCTIConnectorHelper`` instance.
logger: The ``ConnectorLogger`` instance for logging.
"""
self.id = work_id
self.name = work_name
self._helper = helper
self._logger = logger
self._closed = False
self._has_sent_bundles = False

self.logger.debug(f"{self.__class__.__name__} instantiated successfully")

@classmethod
def create(
cls,
helper: OpenCTIConnectorHelper,
logger: ConnectorLogger,
work_name: str,
) -> _Work:
def create(cls, helper: OpenCTIConnectorHelper, work_name: str) -> _Work:
"""Create a new work in OpenCTI and return a ``_Work`` instance.

This classmethod encapsulates the OpenCTI API call to initiate a work,
so that ``WorkManager`` does not need to call the API directly.

Args:
helper: The ``OpenCTIConnectorHelper`` instance.
logger: The ``ConnectorLogger`` instance for logging.
work_name: The name of the work, displayed in the OpenCTI UI.

Returns:
A new ``_Work`` instance wrapping the created work.
"""
work_id: str = helper.api.work.initiate_work(helper.connect_id, work_name)
logger.info(
f"Work '{work_id}' initiated",
{"work_name": work_name},
cls.logger.debug(
"Work created",
{"work_name": work_name, "work_id": work_id},
)
return cls(work_id, work_name, helper, logger)
return cls(helper, work_id, work_name)

def send_bundle(self, bundle_objects: list[Any], **kwargs: Any) -> None:
"""Create a STIX bundle from objects and send it to OpenCTI.
Expand All @@ -96,11 +94,19 @@ def send_bundle(self, bundle_objects: list[Any], **kwargs: Any) -> None:
"""
stix_objects = self._to_stix(bundle_objects)
bundle = self._helper.stix2_create_bundle(stix_objects)
bundles_sent = self._helper.send_stix2_bundle(bundle, work_id=self.id, **kwargs)
bundles_sent = self._helper.send_stix2_bundle(
bundle,
work_id=self.id,
**kwargs,
)
self._has_sent_bundles = True
self._logger.info(
self.logger.info(
"Sent STIX objects to OpenCTI",
{"bundles_sent": str(len(bundles_sent))},
{
"bundles_sent": len(bundles_sent),
"work_name": self.name,
"work_id": self.id,
},
)

def success(self, message: str) -> None:
Expand All @@ -110,7 +116,10 @@ def success(self, message: str) -> None:
message: A completion message stored alongside the work.
"""
self._helper.api.work.to_processed(self.id, message)
self._logger.info(message)
self.logger.debug(
"Work marked as completed on OpenCTI",
{"work_name": self.name, "work_id": self.id, "message": message},
)
self._closed = True

def fail(self, message: str) -> None:
Expand All @@ -120,7 +129,10 @@ def fail(self, message: str) -> None:
message: An error message stored alongside the work.
"""
self._helper.api.work.to_processed(self.id, message, in_error=True)
self._logger.error(message)
self.logger.debug(
"Work marked as failed on OpenCTI",
{"work_name": self.name, "work_id": self.id, "message": message},
)
self._closed = True

def _delete(self) -> None:
Expand All @@ -131,9 +143,9 @@ def _delete(self) -> None:
the ``WorkManager`` to clean up orphaned or invalid works.
"""
self._helper.api.work.delete(id=self.id)
self._logger.info(
"Work deleted",
{"work_id": self.id},
self.logger.debug(
"Work deleted on OpenCTI",
{"work_name": self.name, "work_id": self.id},
)
self._closed = True

Expand Down Expand Up @@ -168,29 +180,28 @@ class WorkManager:

Example::

work_manager = WorkManager(helper, logger)
work_manager = WorkManager(helper)
with work_manager:
work_manager.send(stix_objects, "Import indicators")
work_manager.send(more_objects, "Import indicators") # same work
work_manager.send(reports, "Import reports") # closes previous, new work
# work auto-closed
"""

def __init__(
self,
helper: OpenCTIConnectorHelper,
) -> None:
logger: ClassVar[BaseLogger] = sdk_logger.get_child("WorkManager")

def __init__(self, helper: OpenCTIConnectorHelper) -> None:
"""Initialize the work manager.

Args:
helper: The ``OpenCTIConnectorHelper`` instance.
logger: The ``ConnectorLogger`` instance.
"""
self._helper = helper
self._logger = ConnectorLogger(helper)
self._current_work: _Work | None = None
self._active = False

self.logger.debug(f"{self.__class__.__name__} instantiated successfully")

def __enter__(self) -> WorkManager:
"""Enter the context manager."""
self._current_work = None
Expand All @@ -211,11 +222,34 @@ def __exit__(
"""
if self._current_work is not None and not self._current_work._closed:
if not self._current_work._has_sent_bundles:
self.logger.info(
"Zero bundles were sent, deleting work",
{
"work_name": self._current_work.name,
"work_id": self._current_work.id,
},
)
self._current_work._delete()
elif exc_type is not None:
self._current_work.fail(f"Work failed with error: {exc_val}")
message = f"Work failed with error: {exc_val}"
self.logger.error(
message,
{
"work_name": self._current_work.name,
"work_id": self._current_work.id,
},
)
self._current_work.fail(message)
else:
self._current_work.success("Work completed successfully")
message = "Work completed successfully"
self.logger.info(
message,
{
"work_name": self._current_work.name,
"work_id": self._current_work.id,
},
)
self._current_work.success(message)
self._current_work = None
self._active = False

Expand All @@ -237,19 +271,53 @@ def send(
(e.g. ``cleanup_inconsistent_bundle``, ``update``, ``entities_types``).
"""
if not bundle_objects:
self.logger.info(
"No objects to send",
{"work_name": work_name},
)
return

if not self._active:
msg = "WorkManager.send() must be called inside a 'with' block."
raise RuntimeError(msg)

if self._current_work is None or self._current_work.name != work_name:
self._close_current_work()
self._current_work = _Work.create(self._helper, self._logger, work_name)

self.logger.info(
"Creating a new work",
{
"new_work_name": work_name,
"previous_work_name": (
self._current_work.name
if self._current_work is not None
else None
),
},
)
self._current_work = _Work.create(self._helper, work_name)

self._current_work.send_bundle(bundle_objects, **kwargs)

def _close_current_work(self) -> None:
"""Close the current work if it exists and is not already closed."""
if self._current_work is not None and not self._current_work._closed:
if not self._current_work._has_sent_bundles:
self.logger.info(
"Zero bundles were sent, deleting work",
{
"work_name": self._current_work.name,
"work_id": self._current_work.id,
},
)
self._current_work._delete()
else:
self._current_work.success("Work completed successfully")
message = "Work completed successfully"
self.logger.info(
message,
{
"work_name": self._current_work.name,
"work_id": self._current_work.id,
},
)
self._current_work.success(message)
Loading
Loading