diff --git a/README.md b/README.md index 95293004..0751e884 100644 --- a/README.md +++ b/README.md @@ -258,6 +258,7 @@ When running againts multiple destination organizations, a seperate working dire | roles | Sync Datadog roles. | | rum_applications | Sync Datadog RUM applications. | | rum_metrics | Sync Datadog RUM-based metrics. | +| rum_operations | Sync Datadog RUM operations. | | rum_permanent_retention_filters | Sync Datadog permanent RUM retention filters (configure-only). | | rum_retention_filters | Sync Datadog RUM retention filters (generic + exclusion). | | rum_retention_filters_order | Sync Datadog RUM retention filters order. | @@ -366,6 +367,7 @@ See [Supported resources](#supported-resources) section below for potential reso | roles | - | | rum_applications | - | | rum_metrics | - | +| rum_operations | rum_applications | | rum_permanent_retention_filters | rum_applications | | rum_retention_filters | rum_applications | | rum_retention_filters_order | rum_applications, rum_retention_filters | diff --git a/datadog_sync/model/rum_operations.py b/datadog_sync/model/rum_operations.py new file mode 100644 index 00000000..b5bb15d6 --- /dev/null +++ b/datadog_sync/model/rum_operations.py @@ -0,0 +1,122 @@ +# 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 + +from datadog_sync.utils.base_resource import BaseResource, ResourceConfig +from datadog_sync.utils.custom_client import PaginationConfig +from datadog_sync.utils.resource_utils import CustomClientHTTPError + +if TYPE_CHECKING: + from datadog_sync.utils.custom_client import CustomClient + + +class RUMOperations(BaseResource): + """RUM operations. + + A RUM operation is a static definition (name, display_name, category, query, + journey rules) tied to a RUM application via ``attributes.application_id`` + (a real body field, remapped to the destination app id via + ``resource_connections``). The ``query`` fields in the journey are RUM query + filters (e.g. ``@type:view``), not session ids, so no session-id stripping is + required. List is via the ``/rum/operations/search`` GET endpoint; create is + POST, update is PUT, delete is DELETE. + """ + + resource_type = "rum_operations" + resource_config = ResourceConfig( + base_path="/api/v2/rum/operations", + excluded_attributes=[ + "id", + "attributes.created_at", + "attributes.created_by", + "attributes.updated_at", + "attributes.updated_by", + "attributes.org_id", + ], + resource_connections={ + "rum_applications": ["attributes.application_id"], + }, + skip_resource_mapping=True, + ) + # Additional RUMOperations specific attributes + _search_path = "/api/v2/rum/operations/search" + pagination_config = PaginationConfig( + page_size=100, + page_size_param="page[limit]", + page_number_param="page[offset]", + page_number_func=lambda idx, page_size, page_number: page_number + page_size, + remaining_func=lambda *args: 1, + ) + + async def get_resources(self, client: CustomClient) -> List[Dict]: + resp = await client.paginated_request(client.get)( + self._search_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"] + + 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]: + destination_client = self.config.destination_client + # create data has no id (server-assigned) + resource.pop("id", None) + + # Destination reconciliation: skip_resource_mapping=True means the apply + # pre-pass never lists destination operations, so create_resource + # always runs when state is absent. Before POSTing, search for a + # matching operation at the destination (by name only — operation + # names are unique within an org) and adopt it via update instead of + # creating a duplicate (409 Conflict). + # + # Matching by name only (not name + application_id) because the + # destination app may have been recreated with a new ID since the + # last sync, so the old operation's application_id won't match. + op_name = resource.get("attributes", {}).get("name", "") + try: + existing = await destination_client.get(self._search_path) + for op in existing.get("data", []): + op_attrs = op.get("attributes", {}) + if op_attrs.get("name") == op_name: + self.config.state.destination[self.resource_type][_id] = op + return await self.update_resource(_id, resource) + except CustomClientHTTPError as e: + if e.status_code != 404: + raise + + 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 + destination_id = self.config.state.destination[self.resource_type][_id]["id"] + resource["id"] = destination_id + payload = {"data": resource} + resp = await destination_client.put( + self.resource_config.base_path + f"/{destination_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 2b462628..46ff0476 100644 --- a/datadog_sync/models/__init__.py +++ b/datadog_sync/models/__init__.py @@ -29,6 +29,7 @@ from datadog_sync.model.roles import Roles from datadog_sync.model.rum_applications import RUMApplications from datadog_sync.model.rum_metrics import RUMMetrics +from datadog_sync.model.rum_operations import RUMOperations from datadog_sync.model.rum_permanent_retention_filters import RUMPermanentRetentionFilters from datadog_sync.model.rum_retention_filters import RUMRetentionFilters from datadog_sync.model.rum_retention_filters_order import RUMRetentionFiltersOrder diff --git a/tests/unit/test_rum_operations.py b/tests/unit/test_rum_operations.py new file mode 100644 index 00000000..308d9785 --- /dev/null +++ b/tests/unit/test_rum_operations.py @@ -0,0 +1,182 @@ +# 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 RUMOperations resource model. + +A RUM operation is a static definition (name, display_name, category, query, +journey rules) tied to a RUM application via ``attributes.application_id`` (a +real body field, remapped to the destination app id via ``resource_connections``). +The ``query`` fields in the journey are RUM query filters (e.g. ``@type:view``), +not session ids, so no session-id stripping is required. List is via the +``/rum/operations/search`` GET endpoint. +""" + +import asyncio +from collections import defaultdict +from unittest.mock import AsyncMock, MagicMock + +from datadog_sync.model.rum_operations import RUMOperations + + +def _run(coro): + loop = asyncio.new_event_loop() + try: + return loop.run_until_complete(coro) + finally: + loop.close() + + +def _op(_id, app_id="app-src", name="checkout-flow"): + return { + "id": _id, + "type": "rum_operations", + "attributes": { + "application_id": app_id, + "name": name, + "display_name": "Checkout Flow", + "category": "ux", + "description": "checkout journey", + "tags": ["env:prod"], + "feature_ids": [], + "journey_rum": {"rum_steps": []}, + }, + } + + +def test_get_resources_hits_search_endpoint(): + ops = RUMOperations(MagicMock()) + client = AsyncMock() + # paginated_request(func) returns a wrapper coroutine; mock the wrapper + wrapper_mock = AsyncMock(return_value=[_op("op-1"), _op("op-2")]) + client.paginated_request = MagicMock(return_value=wrapper_mock) + + resources = _run(ops.get_resources(client)) + + assert resources == [_op("op-1"), _op("op-2")] + # Verify paginated_request was called with client.get + client.paginated_request.assert_called_once_with(client.get) + # Verify the wrapper was called with the search path and pagination config + wrapper_mock.assert_called_once_with( + "/api/v2/rum/operations/search", + pagination_config=ops.pagination_config, + ) + + +def test_import_resource_by_id_gets_and_returns(): + ops = RUMOperations(MagicMock()) + source = AsyncMock() + source.get = AsyncMock(return_value={"data": _op("op-1")}) + ops.config.source_client = source + + _id, data = _run(ops.import_resource(_id="op-1")) + + assert _id == "op-1" + source.get.assert_awaited_once_with("/api/v2/rum/operations/op-1") + + +def test_import_resource_passthrough(): + ops = RUMOperations(MagicMock()) + ops.config.source_client = AsyncMock() + resource = _op("op-1") + _id, data = _run(ops.import_resource(resource=resource)) + assert _id == "op-1" + assert data is resource + ops.config.source_client.get.assert_not_awaited() + + +def test_create_resource_posts_without_id(): + ops = RUMOperations(MagicMock()) + dest = AsyncMock() + # Reconciliation GET returns no matching operation + dest.get = AsyncMock(return_value={"data": []}) + dest.post = AsyncMock(return_value={"data": _op("op-dst", app_id="app-dst")}) + ops.config.destination_client = dest + + resource = _op("op-1", app_id="app-dst") + _id, data = _run(ops.create_resource("op-1", resource)) + + assert _id == "op-1" + assert data["id"] == "op-dst" + # create data has no id (server-assigned); it must be popped before POST + assert "id" not in resource + dest.post.assert_awaited_once() + post_url, post_payload = dest.post.await_args.args + assert post_url == "/api/v2/rum/operations" + assert post_payload == {"data": resource} + + +def test_create_resource_reconciles_existing_destination(): + """When a matching operation already exists at the destination (same name), + create_resource hydrates state and delegates to update. Matches by name only + (not name + application_id) because the destination app may have been + recreated with a new ID.""" + ops = RUMOperations(MagicMock()) + dest = AsyncMock() + # Existing op has a DIFFERENT app_id (old app from previous sync) + existing = _op("op-existing", app_id="old-app-id", name="checkout-flow") + dest.get = AsyncMock(return_value={"data": [existing]}) + dest.put = AsyncMock(return_value={"data": existing}) + dest.post = AsyncMock() + ops.config.destination_client = dest + ops.config.state = MagicMock() + ops.config.state.destination = defaultdict(dict) + + # Current sync has a new app_id (newly created destination app) + resource = _op("op-1", app_id="new-app-id", name="checkout-flow") + _id, data = _run(ops.create_resource("op-1", resource)) + + assert _id == "op-1" + dest.post.assert_not_awaited() + dest.put.assert_awaited_once() + # state was hydrated with the existing destination operation + assert ops.config.state.destination["rum_operations"]["op-1"]["id"] == "op-existing" + + +def test_update_resource_puts_destination_id(): + ops = RUMOperations(MagicMock()) + dest = AsyncMock() + dest.put = AsyncMock(return_value={"data": _op("op-dst", app_id="app-dst", name="updated")}) + ops.config.destination_client = dest + ops.config.state = MagicMock() + ops.config.state.destination = defaultdict(dict) + ops.config.state.destination["rum_operations"]["op-1"] = {"id": "op-dst"} + + resource = _op("op-1", app_id="app-dst", name="updated") + _id, data = _run(ops.update_resource("op-1", resource)) + + assert _id == "op-1" + assert resource["id"] == "op-dst" + dest.put.assert_awaited_once() + put_url, put_payload = dest.put.await_args.args + assert put_url == "/api/v2/rum/operations/op-dst" + assert put_payload == {"data": resource} + + +def test_delete_resource_deletes_destination_id(): + ops = RUMOperations(MagicMock()) + dest = AsyncMock() + ops.config.destination_client = dest + ops.config.state = MagicMock() + ops.config.state.destination = defaultdict(dict) + ops.config.state.destination["rum_operations"]["op-1"] = {"id": "op-dst"} + + _run(ops.delete_resource("op-1")) + + dest.delete.assert_awaited_once_with("/api/v2/rum/operations/op-dst") + + +def test_connect_resources_remaps_application_id(): + ops = RUMOperations(MagicMock()) + ops.config.state = MagicMock() + ops.config.state.destination = defaultdict(dict) + ops.config.state.destination["rum_applications"]["app-src"] = {"id": "app-dst"} + ops.config.skip_failed_resource_connections = False + ops.config.logger = MagicMock() + + resource = _op("op-1", app_id="app-src") + ops.connect_resources("op-1", resource) + + assert resource["attributes"]["application_id"] == "app-dst"