From 45fc0aa5353093ccffada2a4bb44ea993b997755 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 25 Sep 2026 16:24:20 -0400 Subject: [PATCH 1/4] feat(rum): add rum_operations resource RUM operations are static definitions (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 -- the feasibility spike confirmed operations are clean static definitions. List is via the /rum/operations/search GET endpoint; create is POST, update is PUT, delete is DELETE. Create data has no id (server-assigned); update sets id to the destination id. - datadog_sync/model/rum_operations.py (new) - datadog_sync/models/__init__.py -- register - tests/unit/test_rum_operations.py (new) -- 7 unit tests - README.md -- add rum_operations (depends on rum_applications) Integration tests + VCR cassettes deferred (require sandbox-org API access). --- README.md | 2 + datadog_sync/model/rum_operations.py | 87 ++++++++++++++++ datadog_sync/models/__init__.py | 1 + tests/unit/test_rum_operations.py | 145 +++++++++++++++++++++++++++ 4 files changed, 235 insertions(+) create mode 100644 datadog_sync/model/rum_operations.py create mode 100644 tests/unit/test_rum_operations.py 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..20f1066c --- /dev/null +++ b/datadog_sync/model/rum_operations.py @@ -0,0 +1,87 @@ +# 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 + +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" + + async def get_resources(self, client: CustomClient) -> List[Dict]: + resp = await client.get(self._search_path) + + return resp["data"] + + 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) + 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..ceef098b --- /dev/null +++ b/tests/unit/test_rum_operations.py @@ -0,0 +1,145 @@ +# 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() + client.get = AsyncMock(return_value={"data": [_op("op-1")], "meta": {}}) + + resources = _run(ops.get_resources(client)) + + assert resources == [_op("op-1")] + client.get.assert_awaited_once_with("/api/v2/rum/operations/search") + + +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() + 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_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" From 110df688378027deac18c8cd8266457e1de3696c Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Fri, 25 Sep 2026 17:24:53 -0400 Subject: [PATCH 2/4] fix(rum): paginate rum_operations search endpoint Per review: /rum/operations/search is paginated but get_resources only returned the first page. Switch to client.paginated_request with PaginationConfig (page[limit]=100, page[offset] incrementing by page_size) so all operations are imported, not just the first page. Test updated to verify paginated_request is used. --- datadog_sync/model/rum_operations.py | 15 +++++++++++++-- tests/unit/test_rum_operations.py | 14 +++++++++++--- 2 files changed, 24 insertions(+), 5 deletions(-) diff --git a/datadog_sync/model/rum_operations.py b/datadog_sync/model/rum_operations.py index 20f1066c..4741e5a2 100644 --- a/datadog_sync/model/rum_operations.py +++ b/datadog_sync/model/rum_operations.py @@ -7,6 +7,7 @@ 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 if TYPE_CHECKING: from datadog_sync.utils.custom_client import CustomClient @@ -42,11 +43,21 @@ class RUMOperations(BaseResource): ) # 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.get(self._search_path) + resp = await client.paginated_request(client.get)( + self._search_path, + pagination_config=self.pagination_config, + ) - return resp["data"] + return resp async def import_resource(self, _id: Optional[str] = None, resource: Optional[Dict] = None) -> Tuple[str, Dict]: if _id: diff --git a/tests/unit/test_rum_operations.py b/tests/unit/test_rum_operations.py index ceef098b..7cdc19c5 100644 --- a/tests/unit/test_rum_operations.py +++ b/tests/unit/test_rum_operations.py @@ -49,12 +49,20 @@ def _op(_id, app_id="app-src", name="checkout-flow"): def test_get_resources_hits_search_endpoint(): ops = RUMOperations(MagicMock()) client = AsyncMock() - client.get = AsyncMock(return_value={"data": [_op("op-1")], "meta": {}}) + # 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")] - client.get.assert_awaited_once_with("/api/v2/rum/operations/search") + 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(): From b4b711a202239066d8552c01b71bf3c4496cb926 Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Thu, 1 Oct 2026 14:13:07 -0400 Subject: [PATCH 3/4] fix(rum): add destination reconciliation to rum_operations create Per integration test: rum_operations has skip_resource_mapping=True, so the apply pre-pass never lists destination operations. On a second run, create_resource POSTs unconditionally and gets 409 Conflict. Fix: create_resource now searches the destination's operations (via the search endpoint) before POSTing. If a matching operation is found (same name + application_id), it hydrates state and delegates to update instead of creating a duplicate. New test verifies the reconciliation path. --- datadog_sync/model/rum_operations.py | 20 ++++++++++++++++++++ tests/unit/test_rum_operations.py | 25 +++++++++++++++++++++++++ 2 files changed, 45 insertions(+) diff --git a/datadog_sync/model/rum_operations.py b/datadog_sync/model/rum_operations.py index 4741e5a2..55fa3047 100644 --- a/datadog_sync/model/rum_operations.py +++ b/datadog_sync/model/rum_operations.py @@ -8,6 +8,7 @@ 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 @@ -76,6 +77,25 @@ 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 + application_id) and + # adopt it via update instead of creating a duplicate (409 Conflict). + op_name = resource.get("attributes", {}).get("name", "") + app_id = resource.get("attributes", {}).get("application_id", "") + 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 and op_attrs.get("application_id") == app_id: + 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"] diff --git a/tests/unit/test_rum_operations.py b/tests/unit/test_rum_operations.py index 7cdc19c5..993e03bd 100644 --- a/tests/unit/test_rum_operations.py +++ b/tests/unit/test_rum_operations.py @@ -90,6 +90,8 @@ def test_import_resource_passthrough(): 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 @@ -106,6 +108,29 @@ def test_create_resource_posts_without_id(): assert post_payload == {"data": resource} +def test_create_resource_reconciles_existing_destination(): + """When a matching operation already exists at the destination (same name + + application_id), create_resource hydrates state and delegates to update.""" + ops = RUMOperations(MagicMock()) + dest = AsyncMock() + existing = _op("op-existing", app_id="app-dst", 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) + + resource = _op("op-1", app_id="app-dst", 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() From dcf0952b7bff6951e0164b29ea28068fb36311ea Mon Sep 17 00:00:00 2001 From: Michael Richey Date: Thu, 1 Oct 2026 14:48:34 -0400 Subject: [PATCH 4/4] fix(rum): match operations by name only in reconciliation Per integration test re-test: the reconciliation in create_resource matched by name AND application_id, but the existing destination operation has an application_id from a previous destination app (created by an earlier sync run). When state is wiped and a new destination app is created, the old operation's application_id doesn't match the new app's ID, so the match fails and the sync POSTs unconditionally (409 Conflict). Fix: match by name only (operation names are unique within an org). Updated test to verify matching works even when the existing op has a different application_id than the current sync. --- datadog_sync/model/rum_operations.py | 12 ++++++++---- tests/unit/test_rum_operations.py | 12 ++++++++---- 2 files changed, 16 insertions(+), 8 deletions(-) diff --git a/datadog_sync/model/rum_operations.py b/datadog_sync/model/rum_operations.py index 55fa3047..b5bb15d6 100644 --- a/datadog_sync/model/rum_operations.py +++ b/datadog_sync/model/rum_operations.py @@ -81,15 +81,19 @@ async def create_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: # 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 + application_id) and - # adopt it via update instead of creating a duplicate (409 Conflict). + # 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", "") - app_id = resource.get("attributes", {}).get("application_id", "") 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 and op_attrs.get("application_id") == app_id: + 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: diff --git a/tests/unit/test_rum_operations.py b/tests/unit/test_rum_operations.py index 993e03bd..308d9785 100644 --- a/tests/unit/test_rum_operations.py +++ b/tests/unit/test_rum_operations.py @@ -109,11 +109,14 @@ def test_create_resource_posts_without_id(): def test_create_resource_reconciles_existing_destination(): - """When a matching operation already exists at the destination (same name + - application_id), create_resource hydrates state and delegates to update.""" + """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("op-existing", app_id="app-dst", name="checkout-flow") + # 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() @@ -121,7 +124,8 @@ def test_create_resource_reconciles_existing_destination(): ops.config.state = MagicMock() ops.config.state.destination = defaultdict(dict) - resource = _op("op-1", app_id="app-dst", name="checkout-flow") + # 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"