Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down Expand Up @@ -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 |
Expand Down
122 changes: 122 additions & 0 deletions datadog_sync/model/rum_operations.py
Original file line number Diff line number Diff line change
@@ -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']}"
)
1 change: 1 addition & 0 deletions datadog_sync/models/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
182 changes: 182 additions & 0 deletions tests/unit/test_rum_operations.py
Original file line number Diff line number Diff line change
@@ -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"
Loading