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
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -253,6 +253,7 @@ When running againts multiple destination organizations, a seperate working dire
| metrics_metadata | Sync Datadog metric metadata. |
| monitors | Sync Datadog monitors. |
| notebooks | Sync Datadog notebooks. |
| observability_pipelines | Sync Datadog Observability Pipelines. |
| powerpacks | Sync Datadog powerpacks. |
| restriction_policies | Sync Datadog restriction policies. |
| roles | Sync Datadog roles. |
Expand Down
95 changes: 95 additions & 0 deletions datadog_sync/model/observability_pipelines.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
# Unless explicitly stated otherwise all files in this repository are licensed
# under the 3-clause BSD style license (see LICENSE).
# This product includes software developed at Datadog (https://www.datadoghq.com/).
# Copyright 2019 Datadog, Inc.

from __future__ import annotations
from typing import TYPE_CHECKING, Optional, List, Dict, Tuple, cast

from datadog_sync.utils.base_resource import BaseResource, ResourceConfig
from datadog_sync.utils.custom_client import PaginationConfig

if TYPE_CHECKING:
from datadog_sync.utils.custom_client import CustomClient


def _op_remaining_func(idx, resp, page_size, page_number):
"""remaining_func for the OP list response.

The OP API returns total count as ``meta.totalCount`` (camelCase), not the
default ``meta.page.total_count`` used by most other v2 endpoints. If ``meta``
or ``totalCount`` is absent, return a negative value so pagination stops
after the current page (the ``resp_len < page_size`` break in
``paginated_request`` already handles the last-page case).
"""
meta = resp.get("meta") or {}
total = meta.get("totalCount", 0)
return total - page_size * (page_number + 1)


class ObservabilityPipelines(BaseResource):
resource_type = "observability_pipelines"
resource_config = ResourceConfig(
base_path="/api/v2/obs-pipelines/pipelines",
excluded_attributes=["id"],
resource_mapping_key="id",
)
# The OP list endpoint paginates with page[size]/page[number] (matching the
# default param names) but returns total count as meta.totalCount, not the
# default meta.page.total_count. Use a custom remaining_func so multi-page
# responses don't raise KeyError on the missing meta.page key.
pagination_config = PaginationConfig(
page_size=100,
page_size_param="page[size]",
page_number_param="page[number]",
remaining_func=_op_remaining_func,
)

async def get_resources(self, client: CustomClient) -> List[Dict]:
resp = await client.paginated_request(client.get)(
self.resource_config.base_path, pagination_config=self.pagination_config
)

return resp

async def import_resource(self, _id: Optional[str] = None, resource: Optional[Dict] = None) -> Tuple[str, Dict]:
if _id:
source_client = self.config.source_client
resource = (await source_client.get(self.resource_config.base_path + f"/{_id}"))["data"]
resource = cast(dict, resource)

return resource["id"], resource

async def pre_resource_action_hook(self, _id, resource: Dict) -> None:
pass

async def pre_apply_hook(self) -> None:
pass

async def create_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]:
if _id in self._existing_resources_map:
self.config.state.destination[self.resource_type][_id] = self._existing_resources_map[_id]
return await self.update_resource(_id, resource)

destination_client = self.config.destination_client
payload = {"data": resource}
resp = await destination_client.post(self.resource_config.base_path, payload)

return _id, resp["data"]

async def update_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]:
destination_client = self.config.destination_client
resource["id"] = self.config.state.destination[self.resource_type][_id]["id"]
payload = {"data": resource}
resp = await destination_client.put(
self.resource_config.base_path + f"/{self.config.state.destination[self.resource_type][_id]['id']}",
payload,
)

return _id, resp["data"]

async def delete_resource(self, _id: str) -> None:
destination_client = self.config.destination_client
await destination_client.delete(
self.resource_config.base_path + f"/{self.config.state.destination[self.resource_type][_id]['id']}"
)
1 change: 1 addition & 0 deletions datadog_sync/models/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
from datadog_sync.model.metrics_metadata import MetricsMetadata
from datadog_sync.model.monitors import Monitors
from datadog_sync.model.notebooks import Notebooks
from datadog_sync.model.observability_pipelines import ObservabilityPipelines
from datadog_sync.model.powerpacks import Powerpacks
from datadog_sync.model.restriction_policies import RestrictionPolicies
from datadog_sync.model.roles import Roles
Expand Down
16 changes: 15 additions & 1 deletion datadog_sync/utils/configuration.py
Original file line number Diff line number Diff line change
Expand Up @@ -252,7 +252,9 @@ def _unwrap_exact_match_pattern(pattern: str) -> str:
return _regex_literal_from_exact_match_body(pattern[1:-1])


_ID_FILE_IMPORT_SUPPORTED_TYPES = frozenset({"monitors", "authn_mappings", "team_memberships", "dashboards"})
_ID_FILE_IMPORT_SUPPORTED_TYPES = frozenset(
{"monitors", "authn_mappings", "team_memberships", "dashboards", "observability_pipelines"}
)
Comment thread
michael-richey marked this conversation as resolved.
"""Resource types eligible for --id-file on the import command.

The import path fans out to per-ID GETs via BaseResource.get_resources_by_ids.
Expand All @@ -271,6 +273,12 @@ def _unwrap_exact_match_pattern(pattern: str) -> str:
import_resource(_id=...) -> queue-handler _import_resource(resource=body)
sequence does exactly one GET per dashboard (no double-fetch). Verified via
tests/unit/test_dashboards_id_file.py.

observability_pipelines: ObservabilityPipelines.import_resource(_id=...) does a
real GET to /api/v2/obs-pipelines/pipelines/{id} and returns the full pipeline
body. The default get_resources_by_ids classifies 404/429/5xx/403 without
aborting, satisfying the disaster-recovery continue-past-errors intent.
Verified via tests/unit/test_observability_pipelines_id_file.py.
"""


Expand All @@ -282,6 +290,7 @@ def _unwrap_exact_match_pattern(pattern: str) -> str:
"host_tags",
"metrics_metadata",
"dashboards",
"observability_pipelines",
}
)
"""Resource types eligible for --id-file on the sync command with --minimize-reads.
Expand All @@ -303,6 +312,11 @@ def _unwrap_exact_match_pattern(pattern: str) -> str:
State.get_by_ids constructs the correct key. Added alongside the import
allowlist entry so the union (_ID_FILE_SUPPORTED_TYPES) accepts dashboards
on both paths by design rather than incidentally via the import set.
- observability_pipelines: state key is the pipeline id. Storage layout:
resources/source/observability_pipelines.<id>.json. ID-derivable, so
State.get_by_ids constructs the correct key. Added alongside the import
allowlist entry for the same reason as dashboards — the union accepts it
on the sync state-load path, so it must be explicitly audited here.

Do NOT widen by config — code-level allowlist forces explicit review.
"""
Expand Down
14 changes: 14 additions & 0 deletions tests/integration/resources/test_observability_pipelines.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
# Unless explicitly stated otherwise all files in this repository are licensed
# under the 3-clause BSD style license (see LICENSE).
# This product includes software developed at Datadog (https://www.datadoghq.com/).

from datadog_sync.models import ObservabilityPipelines
from tests.integration.helpers import BaseResourcesTestClass

import pytest


@pytest.mark.skip(reason="Cannot delete these from destination easily")
class TestObservabilityPipelines(BaseResourcesTestClass):
Comment thread
michael-richey marked this conversation as resolved.
resource_type = ObservabilityPipelines.resource_type
field_to_update = "attributes.name"
Loading
Loading