diff --git a/README.md b/README.md index a06c0a89e..1e7eb59e4 100644 --- a/README.md +++ b/README.md @@ -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. | diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py new file mode 100644 index 000000000..997c68ebf --- /dev/null +++ b/datadog_sync/model/observability_pipelines.py @@ -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", "attributes.config.processors"], + 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=50, + 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']}" + ) diff --git a/datadog_sync/models/__init__.py b/datadog_sync/models/__init__.py index 600ed24ca..0d29f1b69 100644 --- a/datadog_sync/models/__init__.py +++ b/datadog_sync/models/__init__.py @@ -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 diff --git a/datadog_sync/utils/configuration.py b/datadog_sync/utils/configuration.py index fe154bda8..9ad59b8e4 100644 --- a/datadog_sync/utils/configuration.py +++ b/datadog_sync/utils/configuration.py @@ -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"} +) """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. @@ -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. """ @@ -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. @@ -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..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. """ diff --git a/tests/integration/resources/test_observability_pipelines.py b/tests/integration/resources/test_observability_pipelines.py new file mode 100644 index 000000000..330fdf401 --- /dev/null +++ b/tests/integration/resources/test_observability_pipelines.py @@ -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): + resource_type = ObservabilityPipelines.resource_type + field_to_update = "attributes.name" diff --git a/tests/unit/test_observability_pipelines.py b/tests/unit/test_observability_pipelines.py new file mode 100644 index 000000000..782d4757e --- /dev/null +++ b/tests/unit/test_observability_pipelines.py @@ -0,0 +1,268 @@ +# 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. + +"""Unit tests for the observability_pipelines resource model.""" + +import asyncio +from collections import defaultdict +from unittest.mock import AsyncMock, MagicMock + +from datadog_sync.model.observability_pipelines import ObservabilityPipelines +from datadog_sync.utils.configuration import init_resources +from datadog_sync.utils.resource_utils import prep_resource + + +def _make_observability_pipelines() -> ObservabilityPipelines: + config = MagicMock() + config.state = MagicMock() + config.state.destination = defaultdict(dict) + config.state.source = defaultdict(dict) + config.skip_failed_resource_connections = False + config.logger = MagicMock() + return ObservabilityPipelines(config) + + +class TestObservabilityPipelinesRegistration: + def test_resource_type(self): + assert ObservabilityPipelines.resource_type == "observability_pipelines" + + def test_config_contract(self): + rc = ObservabilityPipelines.resource_config + assert rc.base_path == "/api/v2/obs-pipelines/pipelines" + assert rc.resource_mapping_key == "id" + assert rc.excluded_attributes == [ + "root['id']", + "root['attributes']['config']['processors']", + ] + assert rc.skip_resource_mapping is False + + def test_pagination_config_reads_meta_total_count(self): + """The OP list response uses meta.totalCount (camelCase), not the + default meta.page.total_count. Verify the custom remaining_func reads + the right key.""" + pc = ObservabilityPipelines.pagination_config + assert pc.page_size == 50 + assert pc.page_size_param == "page[size]" + assert pc.page_number_param == "page[number]" + assert pc.response_list_accessor == "data" + + # Simulate a full first page with 150 total pipelines. + resp = {"data": ["item"] * 50, "meta": {"totalCount": 75}} + remaining = pc.remaining_func(0, resp, 50, 0) + assert remaining == 25 + + # Second (partial) page: remaining should go negative → loop stops. + resp2 = {"data": ["item"] * 25, "meta": {"totalCount": 75}} + remaining2 = pc.remaining_func(1, resp2, 50, 1) + assert remaining2 == -25 + + def test_registered_in_init_resources(self): + config = MagicMock() + resources = init_resources(config) + assert "observability_pipelines" in resources + assert isinstance(resources["observability_pipelines"], ObservabilityPipelines) + + +class TestObservabilityPipelinesGetResources: + def test_get_resources_uses_paginated_request(self): + op = _make_observability_pipelines() + client = MagicMock() + expected = [{"id": "pipe-src-uuid", "type": "pipelines"}, {"id": "pipe-src-2", "type": "pipelines"}] + inner = AsyncMock(return_value=expected) + client.paginated_request = MagicMock(return_value=inner) + + result = asyncio.run(op.get_resources(client)) + + assert result == expected + client.paginated_request.assert_called_once() + args, _ = client.paginated_request.call_args + assert args[0] == client.get + inner.assert_awaited_once_with( + "/api/v2/obs-pipelines/pipelines", + pagination_config=ObservabilityPipelines.pagination_config, + ) + + def test_get_resources_passes_custom_pagination_config(self): + """get_resources must pass the model's custom pagination_config (with + the OP-specific remaining_func) to paginated_request, not rely on the + client's default.""" + op = _make_observability_pipelines() + client = MagicMock() + inner = AsyncMock(return_value=[]) + client.paginated_request = MagicMock(return_value=inner) + + asyncio.run(op.get_resources(client)) + + _, kwargs = inner.call_args + assert kwargs.get("pagination_config") is ObservabilityPipelines.pagination_config + + +class TestObservabilityPipelinesImportResource: + def test_import_resource_passthrough(self): + op = _make_observability_pipelines() + resource = { + "id": "pipe-src-uuid", + "type": "pipelines", + "attributes": {"name": "test-pipeline", "config": {"sources": [], "destinations": []}}, + } + + _id, result = asyncio.run(op.import_resource(resource=resource)) + + assert _id == "pipe-src-uuid" + assert result == resource + + def test_import_resource_by_id_does_get(self): + op = _make_observability_pipelines() + resource = { + "id": "pipe-src-uuid", + "type": "pipelines", + "attributes": {"name": "test-pipeline", "config": {"sources": [], "destinations": []}}, + } + op.config.source_client = AsyncMock() + op.config.source_client.get.return_value = {"data": resource} + + _id, result = asyncio.run(op.import_resource(_id="pipe-src-uuid")) + + assert _id == "pipe-src-uuid" + assert result == resource + op.config.source_client.get.assert_awaited_once() + call_path = op.config.source_client.get.call_args[0][0] + assert call_path == "/api/v2/obs-pipelines/pipelines/pipe-src-uuid" + + +class TestObservabilityPipelinesCreateResource: + def test_create_when_id_absent_posts(self): + op = _make_observability_pipelines() + op._existing_resources_map = {} + resource = { + "type": "pipelines", + "attributes": {"name": "test-pipeline", "config": {"sources": [], "destinations": []}}, + } + resp_data = {"id": "pipe-dst-uuid", "type": "pipelines", "attributes": resource["attributes"]} + op.config.destination_client = AsyncMock() + op.config.destination_client.post.return_value = {"data": resp_data} + + _id, result = asyncio.run(op.create_resource("pipe-src-uuid", resource)) + + assert _id == "pipe-src-uuid" + assert result == resp_data + op.config.destination_client.post.assert_awaited_once() + args, _ = op.config.destination_client.post.call_args + assert args[0] == "/api/v2/obs-pipelines/pipelines" + assert args[1] == {"data": resource} + assert "id" not in resource + + def test_create_when_id_present_delegates_to_update(self): + op = _make_observability_pipelines() + existing = {"id": "pipe-dst-uuid", "type": "pipelines", "attributes": {"name": "test-pipeline"}} + op._existing_resources_map = {"pipe-src-uuid": existing} + resource = { + "type": "pipelines", + "attributes": {"name": "test-pipeline", "config": {"sources": [], "destinations": []}}, + } + op.config.destination_client = AsyncMock() + op.update_resource = AsyncMock(return_value=("pipe-src-uuid", existing)) + + _id, result = asyncio.run(op.create_resource("pipe-src-uuid", resource)) + + assert _id == "pipe-src-uuid" + assert result == existing + assert op.config.state.destination["observability_pipelines"]["pipe-src-uuid"] == existing + op.update_resource.assert_awaited_once_with("pipe-src-uuid", resource) + + +class TestObservabilityPipelinesUpdateResource: + def test_update_puts_with_destination_id(self): + op = _make_observability_pipelines() + op.config.state.destination["observability_pipelines"]["pipe-src-uuid"] = { + "id": "pipe-dst-uuid", + "type": "pipelines", + "attributes": {"name": "test-pipeline"}, + } + resource = { + "type": "pipelines", + "attributes": {"name": "test-pipeline", "config": {"sources": [], "destinations": []}}, + } + resp_data = {"id": "pipe-dst-uuid", "type": "pipelines", "attributes": resource["attributes"]} + op.config.destination_client = AsyncMock() + op.config.destination_client.put.return_value = {"data": resp_data} + + _id, result = asyncio.run(op.update_resource("pipe-src-uuid", resource)) + + assert _id == "pipe-src-uuid" + assert result == resp_data + assert resource["id"] == "pipe-dst-uuid" + op.config.destination_client.put.assert_awaited_once() + args, _ = op.config.destination_client.put.call_args + assert args[0] == "/api/v2/obs-pipelines/pipelines/pipe-dst-uuid" + assert args[1] == {"data": resource} + + +class TestObservabilityPipelinesDeleteResource: + def test_delete_uses_destination_id(self): + op = _make_observability_pipelines() + op.config.state.destination["observability_pipelines"]["pipe-src-uuid"] = {"id": "pipe-dst-uuid"} + op.config.destination_client = AsyncMock() + + asyncio.run(op.delete_resource("pipe-src-uuid")) + + op.config.destination_client.delete.assert_awaited_once() + call_path = op.config.destination_client.delete.call_args[0][0] + assert call_path == "/api/v2/obs-pipelines/pipelines/pipe-dst-uuid" + + +class TestObservabilityPipelinesPrepResource: + def test_prep_resource_strips_id(self): + resource = { + "id": "pipe-src-uuid", + "type": "pipelines", + "attributes": { + "name": "test-pipeline", + "config": {"sources": [], "destinations": []}, + }, + } + + prep_resource(ObservabilityPipelines.resource_config, resource) + + assert "id" not in resource + assert resource["attributes"]["name"] == "test-pipeline" + assert resource["attributes"]["config"] == {"sources": [], "destinations": []} + + def test_prep_resource_strips_processors(self): + """The OP API rejects write payloads containing both 'processors' and + 'processor_groups'. The source API returns 'processors' as a + read-only field nested inside 'config', so it must be stripped before + sending to the destination.""" + resource = { + "id": "pipe-src-uuid", + "type": "pipelines", + "attributes": { + "name": "test-pipeline", + "config": { + "sources": [], + "destinations": [], + "processors": [{"name": "proc-1", "type": "filter"}], + "processor_groups": [{"name": "pg-1", "processors": []}], + }, + }, + } + + prep_resource(ObservabilityPipelines.resource_config, resource) + + assert "id" not in resource + assert "processors" not in resource["attributes"]["config"] + assert resource["attributes"]["config"]["processor_groups"] == [ + {"name": "pg-1", "processors": []} + ] + + +class TestObservabilityPipelinesHooks: + def test_pre_hooks_are_noops(self): + op = _make_observability_pipelines() + + asyncio.run(op.pre_resource_action_hook("pipe-src-uuid", {})) + result = asyncio.run(op.pre_apply_hook()) + + assert result is None diff --git a/tests/unit/test_observability_pipelines_id_file.py b/tests/unit/test_observability_pipelines_id_file.py new file mode 100644 index 000000000..356c36d61 --- /dev/null +++ b/tests/unit/test_observability_pipelines_id_file.py @@ -0,0 +1,101 @@ +# 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/). + +"""Tests for observability_pipelines support in the --id-file allowlist.""" + +import asyncio +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from datadog_sync.model.observability_pipelines import ObservabilityPipelines +from datadog_sync.utils.configuration import ( + _ID_FILE_IMPORT_SUPPORTED_TYPES, + _ID_FILE_SUPPORTED_TYPES, +) + + +class TestObservabilityPipelinesIDFileSupport: + """Tests for observability_pipelines support in _ID_FILE_SUPPORTED_TYPES.""" + + def test_observability_pipelines_in_id_file_import_supported_types(self): + """'observability_pipelines' must be present in _ID_FILE_IMPORT_SUPPORTED_TYPES + so that `import --id-file=- < {"observability_pipelines": [...]}` is accepted.""" + assert "observability_pipelines" in _ID_FILE_IMPORT_SUPPORTED_TYPES, ( + "observability_pipelines must be in _ID_FILE_IMPORT_SUPPORTED_TYPES. " + "If this fails, id-file import support for observability_pipelines is missing." + ) + + def test_observability_pipelines_in_id_file_supported_types_union(self): + """'observability_pipelines' must also be in the union allowlist consulted by + _parse_id_file (rejects unknown types up-front).""" + assert "observability_pipelines" in _ID_FILE_SUPPORTED_TYPES + + def test_observability_pipelines_in_state_load_supported_types(self): + """Adding observability_pipelines to the import allowlist also adds it to + the union (_ID_FILE_SUPPORTED_TYPES), which _parse_id_file consults for + *both* import and sync --minimize-reads. Since the state-load path + scopes by --resources intersection rather than by + _ID_FILE_STATE_LOAD_SUPPORTED_TYPES, observability_pipelines must be + explicitly in the state-load set too — otherwise it's accepted + incidentally via the union without the ID-derivability verification + the set exists to enforce. The state key is the pipeline id (storage + layout: resources/source/observability_pipelines..json), which is + ID-derivable, so it qualifies.""" + from datadog_sync.utils.configuration import _ID_FILE_STATE_LOAD_SUPPORTED_TYPES + + assert "observability_pipelines" in _ID_FILE_STATE_LOAD_SUPPORTED_TYPES, ( + "observability_pipelines is in _ID_FILE_IMPORT_SUPPORTED_TYPES, so the union " + "accepts it on the sync state-load path too; it must be explicitly " + "in _ID_FILE_STATE_LOAD_SUPPORTED_TYPES (state key is ID-derivable)." + ) + + def test_import_resource_id_does_real_get(self): + """import_resource(_id=...) performs a GET to the OP pipelines API and + returns the body — the per-ID fan-out path used by get_resources_by_ids + on id-file import runs.""" + mock_config = MagicMock() + mock_client = AsyncMock() + body = { + "id": "pipe-src-uuid", + "type": "pipelines", + "attributes": {"name": "test-pipeline", "config": {"sources": [], "destinations": []}}, + } + mock_client.get.return_value = {"data": body} + mock_config.source_client = mock_client + op = ObservabilityPipelines(mock_config) + + _id, resource = asyncio.run(op.import_resource(_id="pipe-src-uuid")) + + mock_client.get.assert_awaited_once() + call_path = mock_client.get.call_args[0][0] + assert ( + call_path == "/api/v2/obs-pipelines/pipelines/pipe-src-uuid" + ), f"import_resource(_id=...) must GET /api/v2/obs-pipelines/pipelines/{{id}}; got {call_path!r}" + assert _id == "pipe-src-uuid" + assert resource == body + + def test_import_resource_id_succeeds_for_valid_uuid(self): + """import_resource(pipeline_uuid) (1-arg) completes without error for a + well-formed UUID.""" + mock_config = MagicMock() + pipeline = { + "id": "423368a4-956a-11ef-b92a-da7ad0900005", + "type": "pipelines", + "attributes": {"name": "test-pipeline", "config": {"sources": [], "destinations": []}}, + } + mock_config.source_client = AsyncMock() + mock_config.source_client.get.return_value = {"data": pipeline} + op = ObservabilityPipelines(mock_config) + _id, _ = asyncio.run(op.import_resource(_id=pipeline["id"])) + assert _id == pipeline["id"] + + def test_import_resource_id_api_error_propagates(self): + """import_resource(uuid) propagates HTTP errors from the upstream GET.""" + mock_config = MagicMock() + mock_config.source_client = AsyncMock() + mock_config.source_client.get.side_effect = Exception("HTTP 404") + op = ObservabilityPipelines(mock_config) + with pytest.raises(Exception, match="404"): + asyncio.run(op.import_resource(_id="nonexistent-uuid"))