From 05b961726342d3b724f9771bf0cc4dc5fa2c9b4b Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 25 Sep 2026 17:01:20 -0400 Subject: [PATCH 01/19] feat: add observability_pipelines resource type Add support for syncing Datadog Observability Pipelines via the v2 /api/v2/obs-pipelines/pipelines API: - New ObservabilityPipelines model (BaseResource subclass) with paginated LIST, GET-by-id import, POST create, PUT update, and DELETE. - Match by id (resource_mapping_key='id'), mirroring logs_metrics. The server-assigned id is excluded from create/update payloads via excluded_attributes=['id']; update_resource re-injects the destination id before the PUT. - Register the model in models/__init__.py so init_resources discovers it. - Add 'observability_pipelines' to _ID_FILE_IMPORT_SUPPORTED_TYPES so import --id-file works (default import_resource(_id=...) does a real GET and get_resources_by_ids classifies 404/429/5xx without aborting). - Unit tests cover config contract, registration, all CRUD methods, prep_resource id-stripping, no-op hooks, and id-file allowlist membership. --- datadog_sync/model/observability_pipelines.py | 68 ++++++ datadog_sync/models/__init__.py | 1 + datadog_sync/utils/configuration.py | 10 +- tests/unit/test_observability_pipelines.py | 201 ++++++++++++++++++ .../test_observability_pipelines_id_file.py | 82 +++++++ 5 files changed, 361 insertions(+), 1 deletion(-) create mode 100644 datadog_sync/model/observability_pipelines.py create mode 100644 tests/unit/test_observability_pipelines.py create mode 100644 tests/unit/test_observability_pipelines_id_file.py diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py new file mode 100644 index 000000000..5ab2d3c2a --- /dev/null +++ b/datadog_sync/model/observability_pipelines.py @@ -0,0 +1,68 @@ +# 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 + +if TYPE_CHECKING: + from datadog_sync.utils.custom_client import CustomClient + + +class ObservabilityPipelines(BaseResource): + resource_type = "observability_pipelines" + resource_config = ResourceConfig( + base_path="/api/v2/obs-pipelines/pipelines", + excluded_attributes=["id"], + resource_mapping_key="id", + ) + + async def get_resources(self, client: CustomClient) -> List[Dict]: + resp = await client.paginated_request(client.get)(self.resource_config.base_path) + + 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..59f1c1f55 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. """ diff --git a/tests/unit/test_observability_pipelines.py b/tests/unit/test_observability_pipelines.py new file mode 100644 index 000000000..738ef994c --- /dev/null +++ b/tests/unit/test_observability_pipelines.py @@ -0,0 +1,201 @@ +# 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']"] + assert rc.skip_resource_mapping is False + + 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") + + +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": []} + + +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..399d17383 --- /dev/null +++ b/tests/unit/test_observability_pipelines_id_file.py @@ -0,0 +1,82 @@ +# 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_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")) From 0f77d7e195fd34211e4fce62f9243153bad3ca4a Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 25 Sep 2026 17:01:24 -0400 Subject: [PATCH 02/19] docs: add observability_pipelines to README and integration test stub MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add observability_pipelines row to the supported-resources table in README.md (alphabetical, between notebooks and powerpacks). - Add integration test stub (skipped — OP pipelines hold downstream routing config, so destructive cleanup against a real org is unsafe for CI). --- README.md | 1 + .../resources/test_observability_pipelines.py | 14 ++++++++++++++ 2 files changed, 15 insertions(+) create mode 100644 tests/integration/resources/test_observability_pipelines.py diff --git a/README.md b/README.md index cf8aed4a4..060603143 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/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" From 171dc9cad63b092cb601e4db263a8c7727890e9f Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 25 Sep 2026 17:10:59 -0400 Subject: [PATCH 03/19] fix: use custom PaginationConfig for OP meta.totalCount response shape The default remaining_func in custom_client.py reads resp["meta"]["page"]["total_count"], but the Observability Pipelines v2 list endpoint returns total count as meta.totalCount (camelCase). For orgs with <100 pipelines the single-page early break hides this, but a full first page would raise KeyError on the missing meta.page key, aborting import/sync for this resource type. Add a custom PaginationConfig with a remaining_func that reads meta.totalCount (with a safe fallback to 0 if meta is absent, so pagination stops gracefully). Update get_resources to pass this config explicitly. Add unit tests for the custom remaining_func (multi-page arithmetic) and for verifying get_resources passes the custom config to paginated_request. --- datadog_sync/model/observability_pipelines.py | 29 +++++++++++++- tests/unit/test_observability_pipelines.py | 39 ++++++++++++++++++- 2 files changed, 66 insertions(+), 2 deletions(-) diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py index 5ab2d3c2a..dad17d499 100644 --- a/datadog_sync/model/observability_pipelines.py +++ b/datadog_sync/model/observability_pipelines.py @@ -7,11 +7,26 @@ 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( @@ -19,9 +34,21 @@ class ObservabilityPipelines(BaseResource): 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) + resp = await client.paginated_request(client.get)( + self.resource_config.base_path, pagination_config=self.pagination_config + ) return resp diff --git a/tests/unit/test_observability_pipelines.py b/tests/unit/test_observability_pipelines.py index 738ef994c..ba1e3a8e6 100644 --- a/tests/unit/test_observability_pipelines.py +++ b/tests/unit/test_observability_pipelines.py @@ -35,6 +35,26 @@ def test_config_contract(self): assert rc.excluded_attributes == ["root['id']"] 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 == 100 + 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"] * 100, "meta": {"totalCount": 150}} + remaining = pc.remaining_func(0, resp, 100, 0) + assert remaining == 50 + + # Second (partial) page: remaining should go negative → loop stops. + resp2 = {"data": ["item"] * 50, "meta": {"totalCount": 150}} + remaining2 = pc.remaining_func(1, resp2, 100, 1) + assert remaining2 == -50 + def test_registered_in_init_resources(self): config = MagicMock() resources = init_resources(config) @@ -56,7 +76,24 @@ def test_get_resources_uses_paginated_request(self): 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") + 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: From 81ca09f9b4e767bdef6757536879dbf6766383d8 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Mon, 28 Sep 2026 10:04:32 -0400 Subject: [PATCH 04/19] fix: add observability_pipelines to _ID_FILE_STATE_LOAD_SUPPORTED_TYPES Adding observability_pipelines to _ID_FILE_IMPORT_SUPPORTED_TYPES also adds it to the union _ID_FILE_SUPPORTED_TYPES, which _parse_id_file consults for both import and sync --minimize-reads --id-file. The established contract (documented in the dashboards precedent at tests/unit/test_dashboards_id_file.py:43-59) requires every ID-derivable type accepted on the state-load path to be explicitly listed in _ID_FILE_STATE_LOAD_SUPPORTED_TYPES, so the ID-derivability is audited rather than accepted incidentally via the union. The OP pipeline state key is the pipeline id (storage layout: resources/source/observability_pipelines..json), which is ID-derivable, so it qualifies. Add it to the state-load set and assert membership in a new unit test. --- datadog_sync/utils/configuration.py | 6 ++++++ .../test_observability_pipelines_id_file.py | 19 +++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/datadog_sync/utils/configuration.py b/datadog_sync/utils/configuration.py index 59f1c1f55..9ad59b8e4 100644 --- a/datadog_sync/utils/configuration.py +++ b/datadog_sync/utils/configuration.py @@ -290,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. @@ -311,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/unit/test_observability_pipelines_id_file.py b/tests/unit/test_observability_pipelines_id_file.py index 399d17383..356c36d61 100644 --- a/tests/unit/test_observability_pipelines_id_file.py +++ b/tests/unit/test_observability_pipelines_id_file.py @@ -32,6 +32,25 @@ def test_observability_pipelines_in_id_file_supported_types_union(self): _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 From 82339d942041c6aa4a7085f4b3219ca084439782 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Mon, 28 Sep 2026 13:13:49 -0400 Subject: [PATCH 05/19] chore: retrigger CI (ubuntu-latest runner stuck) From 218333b235c01a13f3aa431b645562230ce87d09 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Mon, 28 Sep 2026 13:50:23 -0400 Subject: [PATCH 06/19] chore: retrigger CI (test jobs stuck on runner) From 2d44c962ca4123721c60989048ccb601fb24e7b8 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Mon, 28 Sep 2026 16:05:38 -0400 Subject: [PATCH 07/19] chore: retrigger CI (integration test suite timed out) From b2a87178b80808d4d715bd1da895e68cf4868c02 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Tue, 29 Sep 2026 09:43:24 -0400 Subject: [PATCH 08/19] chore: retrigger integration tests (queue cleared) From d5e6cf18fd3d85bfe5826a0357a923a3695f4ac9 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Tue, 29 Sep 2026 11:27:09 -0400 Subject: [PATCH 09/19] chore: retrigger integration tests (previous run was cancelled) From 2237093e4d8dc40eb03af1b5de7a1f310a7592c7 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Wed, 30 Sep 2026 10:06:33 -0400 Subject: [PATCH 10/19] chore: retrigger integration tests (retry) From 6df5a4c6b6bbf2b263b93f8488d111e1ed3204d9 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Wed, 30 Sep 2026 10:48:44 -0400 Subject: [PATCH 11/19] chore: retrigger integration tests (4 flaky failures in downtime_schedules) From f9070330d2128c22491032c9fcf4242f990df7e1 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Wed, 30 Sep 2026 11:25:33 -0400 Subject: [PATCH 12/19] docs: remove observability_pipelines from README resource table --- README.md | 1 - 1 file changed, 1 deletion(-) diff --git a/README.md b/README.md index 1e7eb59e4..a06c0a89e 100644 --- a/README.md +++ b/README.md @@ -253,7 +253,6 @@ 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. | From 20e8c7fd6a5991dd250d613b2cb742ff907bfcf6 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Wed, 30 Sep 2026 11:25:57 -0400 Subject: [PATCH 13/19] chore: remove observability_pipelines model and tests Removing observability_pipelines resource type, model, and all associated unit/integration tests. This PR is focused on fixing the currently failing integration tests in main. --- datadog_sync/model/observability_pipelines.py | 95 ------- .../resources/test_observability_pipelines.py | 14 -- tests/unit/test_observability_pipelines.py | 238 ------------------ .../test_observability_pipelines_id_file.py | 101 -------- 4 files changed, 448 deletions(-) delete mode 100644 datadog_sync/model/observability_pipelines.py delete mode 100644 tests/integration/resources/test_observability_pipelines.py delete mode 100644 tests/unit/test_observability_pipelines.py delete mode 100644 tests/unit/test_observability_pipelines_id_file.py diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py deleted file mode 100644 index dad17d499..000000000 --- a/datadog_sync/model/observability_pipelines.py +++ /dev/null @@ -1,95 +0,0 @@ -# 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']}" - ) diff --git a/tests/integration/resources/test_observability_pipelines.py b/tests/integration/resources/test_observability_pipelines.py deleted file mode 100644 index 330fdf401..000000000 --- a/tests/integration/resources/test_observability_pipelines.py +++ /dev/null @@ -1,14 +0,0 @@ -# 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 deleted file mode 100644 index ba1e3a8e6..000000000 --- a/tests/unit/test_observability_pipelines.py +++ /dev/null @@ -1,238 +0,0 @@ -# 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']"] - 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 == 100 - 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"] * 100, "meta": {"totalCount": 150}} - remaining = pc.remaining_func(0, resp, 100, 0) - assert remaining == 50 - - # Second (partial) page: remaining should go negative → loop stops. - resp2 = {"data": ["item"] * 50, "meta": {"totalCount": 150}} - remaining2 = pc.remaining_func(1, resp2, 100, 1) - assert remaining2 == -50 - - 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": []} - - -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 deleted file mode 100644 index 356c36d61..000000000 --- a/tests/unit/test_observability_pipelines_id_file.py +++ /dev/null @@ -1,101 +0,0 @@ -# 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")) From 42ddc87536c2ddcea44ac60c1a76544397586469 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Wed, 30 Sep 2026 11:27:56 -0400 Subject: [PATCH 14/19] Revert "docs: remove observability_pipelines from README resource table" This reverts commit f9070330d2128c22491032c9fcf4242f990df7e1. --- README.md | 1 + 1 file changed, 1 insertion(+) 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. | From bc295971340a71b0262251236d5f141981b97efb Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Wed, 30 Sep 2026 11:27:56 -0400 Subject: [PATCH 15/19] Revert "chore: remove observability_pipelines model and tests" This reverts commit 20e8c7fd6a5991dd250d613b2cb742ff907bfcf6. --- datadog_sync/model/observability_pipelines.py | 95 +++++++ .../resources/test_observability_pipelines.py | 14 ++ tests/unit/test_observability_pipelines.py | 238 ++++++++++++++++++ .../test_observability_pipelines_id_file.py | 101 ++++++++ 4 files changed, 448 insertions(+) create mode 100644 datadog_sync/model/observability_pipelines.py create mode 100644 tests/integration/resources/test_observability_pipelines.py create mode 100644 tests/unit/test_observability_pipelines.py create mode 100644 tests/unit/test_observability_pipelines_id_file.py diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py new file mode 100644 index 000000000..dad17d499 --- /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"], + 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']}" + ) 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..ba1e3a8e6 --- /dev/null +++ b/tests/unit/test_observability_pipelines.py @@ -0,0 +1,238 @@ +# 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']"] + 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 == 100 + 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"] * 100, "meta": {"totalCount": 150}} + remaining = pc.remaining_func(0, resp, 100, 0) + assert remaining == 50 + + # Second (partial) page: remaining should go negative → loop stops. + resp2 = {"data": ["item"] * 50, "meta": {"totalCount": 150}} + remaining2 = pc.remaining_func(1, resp2, 100, 1) + assert remaining2 == -50 + + 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": []} + + +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")) From 07d3c7cd61387939ce6eb48a6cdd409b6f8f472a Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Wed, 30 Sep 2026 11:30:45 -0400 Subject: [PATCH 16/19] chore: retrigger integration tests (retry - 4 pre-existing downtime_schedules failures) From 5d87f322dceda008a424acaaf1d8af9314950c49 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 2 Oct 2026 14:24:41 -0400 Subject: [PATCH 17/19] fix(observability_pipelines): reduce pagination page_size from 100 to 50 The Observability Pipelines API rejects page[size] > 50 with a 400 Bad Request. Change pagination_config.page_size from 100 to 50 so import succeeds. Update unit tests to match the new page_size. --- datadog_sync/model/observability_pipelines.py | 2 +- tests/unit/test_observability_pipelines.py | 14 +++++++------- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py index dad17d499..fb18f3198 100644 --- a/datadog_sync/model/observability_pipelines.py +++ b/datadog_sync/model/observability_pipelines.py @@ -39,7 +39,7 @@ class ObservabilityPipelines(BaseResource): # 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=50, page_size_param="page[size]", page_number_param="page[number]", remaining_func=_op_remaining_func, diff --git a/tests/unit/test_observability_pipelines.py b/tests/unit/test_observability_pipelines.py index ba1e3a8e6..f7957ec0c 100644 --- a/tests/unit/test_observability_pipelines.py +++ b/tests/unit/test_observability_pipelines.py @@ -40,20 +40,20 @@ def test_pagination_config_reads_meta_total_count(self): default meta.page.total_count. Verify the custom remaining_func reads the right key.""" pc = ObservabilityPipelines.pagination_config - assert pc.page_size == 100 + 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"] * 100, "meta": {"totalCount": 150}} - remaining = pc.remaining_func(0, resp, 100, 0) - assert remaining == 50 + 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"] * 50, "meta": {"totalCount": 150}} - remaining2 = pc.remaining_func(1, resp2, 100, 1) - assert remaining2 == -50 + 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() From 951658e1c2891595dc836fb9eee6d83dcf52dff7 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 2 Oct 2026 15:22:06 -0400 Subject: [PATCH 18/19] fix(observability_pipelines): exclude attributes.processors from write payload The OP API rejects create/update payloads containing both 'processors' and 'processor_groups'. The source API returns 'processors' as a read-only field, so it must be stripped via excluded_attributes before sending to the destination. Add 'attributes.processors' to excluded_attributes and a unit test verifying the strip. --- datadog_sync/model/observability_pipelines.py | 2 +- tests/unit/test_observability_pipelines.py | 31 ++++++++++++++++++- 2 files changed, 31 insertions(+), 2 deletions(-) diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py index fb18f3198..58d84a0b2 100644 --- a/datadog_sync/model/observability_pipelines.py +++ b/datadog_sync/model/observability_pipelines.py @@ -31,7 +31,7 @@ class ObservabilityPipelines(BaseResource): resource_type = "observability_pipelines" resource_config = ResourceConfig( base_path="/api/v2/obs-pipelines/pipelines", - excluded_attributes=["id"], + excluded_attributes=["id", "attributes.processors"], resource_mapping_key="id", ) # The OP list endpoint paginates with page[size]/page[number] (matching the diff --git a/tests/unit/test_observability_pipelines.py b/tests/unit/test_observability_pipelines.py index f7957ec0c..6c11c8dd0 100644 --- a/tests/unit/test_observability_pipelines.py +++ b/tests/unit/test_observability_pipelines.py @@ -32,7 +32,10 @@ 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']"] + assert rc.excluded_attributes == [ + "root['id']", + "root['attributes']['processors']", + ] assert rc.skip_resource_mapping is False def test_pagination_config_reads_meta_total_count(self): @@ -227,6 +230,32 @@ def test_prep_resource_strips_id(self): 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' (read-only), + so it must be stripped before sending to the destination.""" + resource = { + "id": "pipe-src-uuid", + "type": "pipelines", + "attributes": { + "name": "test-pipeline", + "processors": [{"name": "proc-1", "type": "filter"}], + "config": { + "sources": [], + "destinations": [], + "processor_groups": [{"name": "pg-1", "processors": []}], + }, + }, + } + + prep_resource(ObservabilityPipelines.resource_config, resource) + + assert "id" not in resource + assert "processors" not in resource["attributes"] + assert resource["attributes"]["config"]["processor_groups"] == [ + {"name": "pg-1", "processors": []} + ] + class TestObservabilityPipelinesHooks: def test_pre_hooks_are_noops(self): From 2599b85e282a22e9d208b883bb357c607acb7301 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 2 Oct 2026 16:07:08 -0400 Subject: [PATCH 19/19] fix(observability_pipelines): correct excluded_attributes path to attributes.config.processors The 'processors' field is nested inside 'config', not directly in 'attributes'. The previous exclusion 'attributes.processors' did not match the actual path 'attributes.config.processors', so the field was not stripped and the API still rejected the payload with 400 Bad Request ('cannot specify both processors and processor_groups'). Fix: change excluded_attributes from 'attributes.processors' to 'attributes.config.processors'. Update unit tests to match. --- datadog_sync/model/observability_pipelines.py | 2 +- tests/unit/test_observability_pipelines.py | 11 ++++++----- 2 files changed, 7 insertions(+), 6 deletions(-) diff --git a/datadog_sync/model/observability_pipelines.py b/datadog_sync/model/observability_pipelines.py index 58d84a0b2..997c68ebf 100644 --- a/datadog_sync/model/observability_pipelines.py +++ b/datadog_sync/model/observability_pipelines.py @@ -31,7 +31,7 @@ class ObservabilityPipelines(BaseResource): resource_type = "observability_pipelines" resource_config = ResourceConfig( base_path="/api/v2/obs-pipelines/pipelines", - excluded_attributes=["id", "attributes.processors"], + excluded_attributes=["id", "attributes.config.processors"], resource_mapping_key="id", ) # The OP list endpoint paginates with page[size]/page[number] (matching the diff --git a/tests/unit/test_observability_pipelines.py b/tests/unit/test_observability_pipelines.py index 6c11c8dd0..782d4757e 100644 --- a/tests/unit/test_observability_pipelines.py +++ b/tests/unit/test_observability_pipelines.py @@ -34,7 +34,7 @@ def test_config_contract(self): assert rc.resource_mapping_key == "id" assert rc.excluded_attributes == [ "root['id']", - "root['attributes']['processors']", + "root['attributes']['config']['processors']", ] assert rc.skip_resource_mapping is False @@ -232,17 +232,18 @@ def test_prep_resource_strips_id(self): def test_prep_resource_strips_processors(self): """The OP API rejects write payloads containing both 'processors' and - 'processor_groups'. The source API returns 'processors' (read-only), - so it must be stripped before sending to the destination.""" + '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", - "processors": [{"name": "proc-1", "type": "filter"}], "config": { "sources": [], "destinations": [], + "processors": [{"name": "proc-1", "type": "filter"}], "processor_groups": [{"name": "pg-1", "processors": []}], }, }, @@ -251,7 +252,7 @@ def test_prep_resource_strips_processors(self): prep_resource(ObservabilityPipelines.resource_config, resource) assert "id" not in resource - assert "processors" not in resource["attributes"] + assert "processors" not in resource["attributes"]["config"] assert resource["attributes"]["config"]["processor_groups"] == [ {"name": "pg-1", "processors": []} ]