Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
05b9617
feat: add observability_pipelines resource type
michael-richey Sep 25, 2026
0f77d7e
docs: add observability_pipelines to README and integration test stub
michael-richey Sep 25, 2026
171dc9c
fix: use custom PaginationConfig for OP meta.totalCount response shape
michael-richey Sep 25, 2026
81ca09f
fix: add observability_pipelines to _ID_FILE_STATE_LOAD_SUPPORTED_TYPES
michael-richey Sep 28, 2026
82339d9
chore: retrigger CI (ubuntu-latest runner stuck)
michael-richey Sep 28, 2026
218333b
chore: retrigger CI (test jobs stuck on runner)
michael-richey Sep 28, 2026
2d44c96
chore: retrigger CI (integration test suite timed out)
michael-richey Sep 28, 2026
b2a8717
chore: retrigger integration tests (queue cleared)
michael-richey Sep 29, 2026
90a3c83
Merge remote-tracking branch 'origin/main' into michael.richey/add-ob…
michael-richey Sep 29, 2026
d5e6cf1
chore: retrigger integration tests (previous run was cancelled)
michael-richey Sep 29, 2026
2237093
chore: retrigger integration tests (retry)
michael-richey Sep 30, 2026
6df5a4c
chore: retrigger integration tests (4 flaky failures in downtime_sche…
michael-richey Sep 30, 2026
f907033
docs: remove observability_pipelines from README resource table
michael-richey Sep 30, 2026
20e8c7f
chore: remove observability_pipelines model and tests
michael-richey Sep 30, 2026
42ddc87
Revert "docs: remove observability_pipelines from README resource table"
michael-richey Sep 30, 2026
bc29597
Revert "chore: remove observability_pipelines model and tests"
michael-richey Sep 30, 2026
07d3c7c
chore: retrigger integration tests (retry - 4 pre-existing downtime_s…
michael-richey Sep 30, 2026
95a6430
Merge remote-tracking branch 'origin/main' into michael.richey/add-ob…
michael-richey Sep 30, 2026
5d87f32
fix(observability_pipelines): reduce pagination page_size from 100 to 50
michael-richey Oct 2, 2026
951658e
fix(observability_pipelines): exclude attributes.processors from writ…
michael-richey Oct 2, 2026
2599b85
fix(observability_pipelines): correct excluded_attributes path to att…
michael-richey Oct 2, 2026
File filter

Filter by extension

Filter by extension

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

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

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

if TYPE_CHECKING:
from datadog_sync.utils.custom_client import CustomClient


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

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


class ObservabilityPipelines(BaseResource):
resource_type = "observability_pipelines"
resource_config = ResourceConfig(
base_path="/api/v2/obs-pipelines/pipelines",
excluded_attributes=["id", "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']}"
)
1 change: 1 addition & 0 deletions datadog_sync/models/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
from datadog_sync.model.metrics_metadata import MetricsMetadata
from datadog_sync.model.monitors import Monitors
from datadog_sync.model.notebooks import Notebooks
from datadog_sync.model.observability_pipelines import ObservabilityPipelines
from datadog_sync.model.powerpacks import Powerpacks
from datadog_sync.model.restriction_policies import RestrictionPolicies
from datadog_sync.model.roles import Roles
Expand Down
16 changes: 15 additions & 1 deletion datadog_sync/utils/configuration.py
Original file line number Diff line number Diff line change
Expand Up @@ -252,7 +252,9 @@ def _unwrap_exact_match_pattern(pattern: str) -> str:
return _regex_literal_from_exact_match_body(pattern[1:-1])


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

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

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


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

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

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

import pytest


@pytest.mark.skip(reason="Cannot delete these from destination easily")
class TestObservabilityPipelines(BaseResourcesTestClass):
resource_type = ObservabilityPipelines.resource_type
field_to_update = "attributes.name"
Loading
Loading