From 803dcf2f62ff4070d3b768235fc5f2741363b5d3 Mon Sep 17 00:00:00 2001 From: michael-richey <41595765+michael-richey@users.noreply.github.com> Date: Fri, 17 Jul 2026 14:48:43 -0400 Subject: [PATCH 1/7] feat(users): key users by handle and add --use-v1-user-api MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Users were identified and linked by email, but the destination enforces user uniqueness on the handle (email is not unique). Two source users that share an email but have distinct handles collapsed onto one derived handle — one user was mislabeled and the other could not be created. - Identify/link users by their exact-case handle: switch resource_mapping_key to attributes.handle. Keep handle out of the read-only v2 create/update payload (pop before POST/PATCH) and out of update diffs (exclude_regex_paths), while retaining it for mapping and v1 creation. - Add an opt-in --use-v1-user-api flag (default off; env DD_USE_V1_USER_API). When set, create users via POST /api/v1/user, which accepts an explicit handle, so every user keeps its own handle and no create collides. The v1 response is the legacy shape, so re-resolve the v2 UUID by exact-case handle (small sleep-free re-query loop for read-after-write visibility); failures re-raise so the apply loop continues (DR-safe). With the flag off, creation is unchanged (v2). - Document the flag in the README; add pytest-cov to test extras. Co-Authored-By: Claude Opus 4.8 --- README.md | 8 + datadog_sync/commands/shared/options.py | 12 + datadog_sync/constants.py | 1 + datadog_sync/model/users.py | 68 ++++- datadog_sync/utils/configuration.py | 3 + setup.cfg | 1 + tests/unit/test_users.py | 334 ++++++++++++++++++++++++ 7 files changed, 423 insertions(+), 4 deletions(-) create mode 100644 tests/unit/test_users.py diff --git a/README.md b/README.md index 6f63cf0e1..28872f6f9 100644 --- a/README.md +++ b/README.md @@ -19,6 +19,7 @@ Datadog cli tool to sync resources across organizations. - [Config file](#config-file) - [Cleanup flag](#cleanup-flag) - [Verify DDR Status Flag](#verify-ddr-status-flag) + - [Create users via the v1 API](#create-users-via-the-v1-api) - [State Files](#state-files) - [Supported resources](#supported-resources) - [Best practices](#best-practices) @@ -210,6 +211,13 @@ For example, `ResourceA` and `ResourceB` are imported and synced, followed by de By default all commands check the Datadog Disaster Recovery (DDR) status of both the source and destination organizations before running. This behavior is controlled by the boolean flag `--verify-ddr-status` or the environment variable `DD_VERIFY_DDR_STATUS`. +#### Create users via the v1 API + +The destination enforces user uniqueness on the `handle`, not the email — multiple users may share an email address. The v2 user create endpoint (`POST /api/v2/users`) does not accept a `handle` (it derives one from the email), so users that share an email collapse onto a single handle: the first create takes a handle that may belong to another user, and later creates for that email return an HTTP 409. + +Passing `--use-v1-user-api` (or setting `DD_USE_V1_USER_API=true`) makes `sync` create users via the v1 user endpoint (`POST /api/v1/user`), which accepts an explicit handle, so every user keeps its own handle and no create collides. The flag is off by default. + + #### Running behind an HTTP proxy By default the tool's HTTP client ignores the environment and talks to Datadog directly. To run it behind a proxy, set `--http-client-trust-env true` (or the environment variable `DD_HTTP_CLIENT_TRUST_ENV=true`). When enabled, the underlying HTTP client honors the standard `HTTP_PROXY`, `HTTPS_PROXY`, and `NO_PROXY` environment variables, as well as credentials from `.netrc`. This option is off by default. Note that when enabled, the configured proxy can observe all Datadog API traffic — including the `DD-API-KEY`, `DD-APPLICATION-KEY`, or JWT headers if it terminates TLS — and `.netrc` credentials may be automatically attached for matching hosts, so only enable this for a proxy you trust. diff --git a/datadog_sync/commands/shared/options.py b/datadog_sync/commands/shared/options.py index a519b6ec7..1de6ea08a 100644 --- a/datadog_sync/commands/shared/options.py +++ b/datadog_sync/commands/shared/options.py @@ -591,6 +591,18 @@ def click_config_file_provider(ctx: Context, opts: CustomOptionClass, value: Non help="Allow self-lockout when syncing restriction policies.", cls=CustomOptionClass, ), + option( + "--use-v1-user-api", + required=False, + envvar=constants.DD_USE_V1_USER_API, + type=bool, + default=False, + show_default=True, + help="Create users via the v1 API (/api/v1/user) instead of v2. v2 cannot set a " + "handle (it derives one from the email), which collides when users share an email; " + "v1 accepts an explicit handle so each user keeps its own. Off by default.", + cls=CustomOptionClass, + ), ] diff --git a/datadog_sync/constants.py b/datadog_sync/constants.py index d62d84854..6c8fefd50 100644 --- a/datadog_sync/constants.py +++ b/datadog_sync/constants.py @@ -29,6 +29,7 @@ DD_SHOW_PROGRESS_BAR = "DD_SHOW_PROGRESS_BAR" DD_VERIFY_SSL_CERTIFICATES = "DD_VERIFY_SSL_CERTIFICATES" DD_ALLOW_PARTIAL_PERMISSIONS_ROLES = "DD_ALLOW_PARTIAL_PERMISSIONS_ROLES" +DD_USE_V1_USER_API = "DD_USE_V1_USER_API" DD_SYNC_JSON = "DD_SYNC_JSON" DD_DATADOG_HOST_OVERRIDE = "DD_DATADOG_HOST_OVERRIDE" diff --git a/datadog_sync/model/users.py b/datadog_sync/model/users.py index 5b4e5a795..52403428e 100644 --- a/datadog_sync/model/users.py +++ b/datadog_sync/model/users.py @@ -30,7 +30,11 @@ class Users(BaseResource): "attributes.status", "attributes.verified", "attributes.service_account", - "attributes.handle", + # NOTE: attributes.handle is deliberately NOT excluded here. It is the + # user mapping key (resource_mapping_key below) and the payload for the + # v1 user creation, so it must survive prep_resource. It is popped + # manually before the v2 POST/PATCH (v2 treats handle as read-only) and + # kept out of update diffs via deep_diff_config.exclude_regex_paths below. "attributes.icon", "attributes.modified_at", "attributes.mfa_enabled", @@ -40,7 +44,11 @@ class Users(BaseResource): "relationships.org", "relationships.team_roles", ], - resource_mapping_key="attributes.email", + resource_mapping_key="attributes.handle", + # Handle is read-only in v2 and is popped from the create/update payload; + # exclude it from update diffs so a source-vs-destination handle difference + # never drives a spurious PATCH of a field v2 will not accept. + deep_diff_config={"ignore_order": True, "exclude_regex_paths": [r".*\['handle'\]"]}, ) # Additional Users specific attributes pagination_config = PaginationConfig( @@ -80,12 +88,63 @@ async def create_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: self.config.state.destination[self.resource_type][_id] = self._existing_resources_map[key] return await self.update_resource(_id, resource) + attributes = resource["attributes"] + if self.config.use_v1_user_api: + # v2 create cannot set a handle (it derives one from the email), which + # collapses distinct-handle users that share an email onto one handle. + # v1 accepts an explicit handle, so each user keeps its own. + return _id, await self._create_via_v1( + attributes.get("handle"), attributes.get("name"), attributes.get("email") + ) + destination_client = self.config.destination_client - resource["attributes"].pop("disabled", None) + # handle is read-only in v2 (derived from email) and must not be sent. + attributes.pop("disabled", None) + attributes.pop("handle", None) resp = await destination_client.post(self.resource_config.base_path, {"data": resource}) - return _id, resp["data"] + async def _create_via_v1(self, handle: Optional[str], name: Optional[str], email: Optional[str]) -> Dict: + """Create the user via the v1 API, which accepts an explicit handle. + + v2 ``POST /api/v2/users`` cannot set a handle (it is derived from the + email), so users that share an email collapse onto one handle and later + creates 409 — and the v2 "winner" is created with the wrong handle. v1 + ``POST /api/v1/user`` accepts a distinct handle. The v1 response is the + legacy user shape, so re-resolve the v2 UUID by handle and return that + record — state and downstream references (roles, team_memberships) key + on the v2 UUID. + """ + if not handle: + raise ValueError("v1 user creation requires a handle") + + destination_client = self.config.destination_client + await destination_client.post("/api/v1/user", {"handle": handle, "name": name, "email": email}) + + user = await self._get_destination_user_by_handle(handle) + if user is None: + raise ValueError("v1-created user not found by handle after create") + return user + + async def _get_destination_user_by_handle(self, handle: str) -> Optional[Dict]: + """Return the destination user whose handle matches exactly, or None. + + Transient HTTP errors are already retried by the client's + ``request_with_retry``; this adds a small, sleep-free re-query loop only + to absorb read-after-write visibility lag after a v1 create. + """ + destination_client = self.config.destination_client + for _ in range(3): + resp = await destination_client.paginated_request(destination_client.get)( + self.resource_config.base_path, + pagination_config=self.pagination_config, + params={"filter": handle}, + ) + for user in resp: + if user.get("attributes", {}).get("handle") == handle: + return user + return None + async def update_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: destination_client = self.config.destination_client @@ -94,6 +153,7 @@ async def update_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: await self.update_user_roles(self.config.state.destination[self.resource_type][_id]["id"], diff) resource["id"] = self.config.state.destination[self.resource_type][_id]["id"] resource.pop("relationships", None) + resource["attributes"].pop("handle", None) resp = await destination_client.patch( self.resource_config.base_path + f"/{self.config.state.destination[self.resource_type][_id]['id']}", {"data": resource}, diff --git a/datadog_sync/utils/configuration.py b/datadog_sync/utils/configuration.py index 755897606..41b391286 100644 --- a/datadog_sync/utils/configuration.py +++ b/datadog_sync/utils/configuration.py @@ -109,6 +109,7 @@ class Configuration(object): max_workers_per_type: Dict[str, int] = field(default_factory=dict) command: str = "" allow_partial_permissions_roles: List[str] = field(default_factory=list) + use_v1_user_api: bool = False resources: Dict[str, BaseResource] = field(default_factory=dict) resources_arg: List[str] = field(default_factory=list) # --id-file: id-targeted import via stdin or file payload. @@ -435,6 +436,7 @@ def build_config(cmd: Command, **kwargs: Optional[Any]) -> Configuration: max_workers_per_type = _parse_max_workers_per_type(max_workers_per_type_raw, known_resource_types) create_global_downtime = kwargs.get("create_global_downtime") validate = kwargs.get("validate") + use_v1_user_api = kwargs.get("use_v1_user_api") or False verify_ddr_status = kwargs.get("verify_ddr_status") backup_before_reset = not kwargs.get("do_not_backup") show_progress_bar = kwargs.get("show_progress_bar") @@ -691,6 +693,7 @@ def build_config(cmd: Command, **kwargs: Optional[Any]) -> Configuration: emit_json=emit_json, command=cmd.value, allow_partial_permissions_roles=allow_partial_permissions_roles, + use_v1_user_api=use_v1_user_api, id_payload=id_payload, max_concurrent_reads=max_concurrent_reads, transient_failure_threshold_pct=transient_failure_threshold_pct, diff --git a/setup.cfg b/setup.cfg index 0b06fe831..e9bc0462e 100644 --- a/setup.cfg +++ b/setup.cfg @@ -56,6 +56,7 @@ tests = black==24.3.0 pytest==8.1.1 pytest-black + pytest-cov pytest-console-scripts pytest-recording==0.13.2 pytest-retry==1.7.0 diff --git a/tests/unit/test_users.py b/tests/unit/test_users.py new file mode 100644 index 000000000..8a74a1db8 --- /dev/null +++ b/tests/unit/test_users.py @@ -0,0 +1,334 @@ +# 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. + +"""Tests for handle-keyed users + the --use-v1-user-api flag. + +The Datadog destination enforces user uniqueness on the ``handle``, not the +``email`` — multiple users may share an email. Keying ``_existing_resources_map`` +by email collapsed distinct-handle users onto one derived handle and caused a +409 Conflict on the second create. These tests cover switching the user mapping +key to the exact-case handle (PR1) and wiring the opt-in --use-v1-user-api flag. + +Tests ``a``/``a2``/``a3``/``f`` are red against ``resource_mapping_key= +"attributes.email"`` and green after the switch to ``"attributes.handle"`` plus +the manual handle pop before the v2 POST. Test ``g`` is a green/green guard that +handle stays excluded from update diffs across the excluded_attributes -> +exclude_regex_paths migration. ``p1``/``p2``/``p3`` guard the flag wiring. +""" + +import asyncio +from unittest.mock import AsyncMock, MagicMock + +import pytest +from click.testing import CliRunner + +from datadog_sync.cli import cli +from datadog_sync.constants import Command +from datadog_sync.model.users import Users +from datadog_sync.utils.configuration import build_config +from datadog_sync.utils.resource_utils import CustomClientHTTPError, check_diff + + +def _http_error(status): + """Build a CustomClientHTTPError with the given status code.""" + response = MagicMock() + response.status = status + response.message = "Conflict" if status == 409 else "Error" + return CustomClientHTTPError(response) + + +def _make_user(handle, email, user_id, name="User"): + """Build a user dict shaped like the v2 GET/create response.""" + return { + "id": user_id, + "type": "users", + "attributes": { + "handle": handle, + "email": email, + "name": name, + "disabled": False, + }, + "relationships": {"roles": {"data": []}}, + } + + +def _base_kwargs(tmp_path): + """Minimal build_config kwargs that avoid network/validation.""" + return dict( + resources="users", + resource_per_file=True, + source_api_key="k", + source_app_key="k", + destination_api_key="k", + destination_app_key="k", + source_api_url="https://example.com", + destination_api_url="https://example.com", + storage_type="local", + source_resources_path=str(tmp_path / "source"), + destination_resources_path=str(tmp_path / "dest"), + max_workers=1, + send_metrics=False, + verify_ddr_status=False, + validate=False, + show_progress_bar=False, + allow_self_lockout=False, + force_missing_dependencies=False, + skip_failed_resource_connections=False, + ) + + +class TestHandleMappingKey: + def test_mapping_key_is_exact_case_handle(self, mock_config): + """a: mapping key is the handle, preserved exact-case (no lowercasing).""" + instance = Users(mock_config) + resource = {"attributes": {"handle": "User-A@example.com", "email": "shared@example.com"}} + assert instance.get_resource_mapping_key(resource) == "User-A@example.com" + + def test_map_keeps_shared_email_distinct_by_handle(self, mock_config): + """a2: two destination users sharing an email but with distinct handles + remain two entries in the map (email keying would collapse to one).""" + instance = Users(mock_config) + dest = [ + _make_user("user-a@example.com", "shared@example.com", "dest-a"), + _make_user("user-b@example.com", "shared@example.com", "dest-b"), + ] + instance.get_resources = AsyncMock(return_value=dest) + asyncio.run(instance.map_existing_resources()) + assert set(instance._existing_resources_map.keys()) == { + "user-a@example.com", + "user-b@example.com", + } + + def test_source_matched_by_handle_not_email(self, mock_config): + """a3: against a pre-existing destination user, a source user is matched + by handle — a same-email/different-handle source is NOT a match.""" + instance = Users(mock_config) + instance.get_resources = AsyncMock( + return_value=[_make_user("user-a@example.com", "shared@example.com", "dest-a")] + ) + asyncio.run(instance.map_existing_resources()) + + same_handle = _make_user("user-a@example.com", "shared@example.com", "src-a") + assert instance.get_resource_mapping_key(same_handle) in instance._existing_resources_map + + diff_handle_same_email = _make_user("user-b@example.com", "shared@example.com", "src-b") + assert instance.get_resource_mapping_key(diff_handle_same_email) not in instance._existing_resources_map + + +class TestV2CreatePayload: + def test_v2_post_body_excludes_handle_and_disabled(self, mock_config): + """f: the v2 create body must not carry read-only handle or disabled.""" + mock_config.use_v1_user_api = False + instance = Users(mock_config) + instance._existing_resources_map = {} + mock_config.destination_client.post = AsyncMock( + return_value={"data": {"id": "dest-x", "attributes": {}}} + ) + src = _make_user("user-a@example.com", "shared@example.com", "src-a") + + asyncio.run(instance.create_resource("src-a", src)) + + mock_config.destination_client.post.assert_called_once() + _, body = mock_config.destination_client.post.call_args.args + attrs = body["data"]["attributes"] + assert "handle" not in attrs + assert "disabled" not in attrs + + def test_v2_patch_body_excludes_handle(self, mock_config): + """f2: the v2 update (PATCH) body must not carry the read-only handle.""" + instance = Users(mock_config) + _id = "src-a" + mock_config.state.destination["users"][_id] = _make_user( + "user-a@example.com", "shared@example.com", "dest-a" + ) + # A differing name forces a diff -> the PATCH branch. + src = _make_user("user-a@example.com", "shared@example.com", "dest-a", name="New Name") + mock_config.destination_client.patch = AsyncMock( + return_value={"data": {"id": "dest-a", "attributes": {}}} + ) + + asyncio.run(instance.update_resource(_id, src)) + + mock_config.destination_client.patch.assert_called_once() + _, body = mock_config.destination_client.patch.call_args.args + assert "handle" not in body["data"]["attributes"] + + +class TestHandleDiffExclusion: + def test_handle_excluded_from_update_diff(self): + """g: two users differing only by handle produce no diff (guards that + handle stays diff-excluded after moving off excluded_attributes).""" + dest = _make_user("user-a@example.com", "shared@example.com", "same-id") + src = _make_user("user-b@example.com", "shared@example.com", "same-id") + assert not check_diff(Users.resource_config, dest, src) + + +class TestV1UserApiFlagWiring: + def test_build_config_flag_true(self, tmp_path): + """p1: --use-v1-user-api flows from kwargs into Configuration.""" + cfg = build_config(Command.IMPORT, use_v1_user_api=True, **_base_kwargs(tmp_path)) + assert cfg.use_v1_user_api is True + + def test_build_config_flag_default_false(self, tmp_path): + """p1: absent flag defaults to False (guards a typo'd kwarg key).""" + cfg = build_config(Command.IMPORT, **_base_kwargs(tmp_path)) + assert cfg.use_v1_user_api is False + + def test_sync_accepts_v1_flag(self): + """p2: the sync command recognizes the flag (exit 2 == usage error).""" + result = CliRunner(mix_stderr=False).invoke( + cli, ["sync", "--use-v1-user-api=true", "--validate=false"] + ) + assert result.exit_code != 2 + + def test_migrate_accepts_v1_flag(self): + """p3: migrate reuses @sync_options, so it recognizes the flag too.""" + result = CliRunner(mix_stderr=False).invoke( + cli, ["migrate", "--use-v1-user-api=true", "--validate=false"] + ) + assert result.exit_code != 2 + + +def _mock_paginated(config, pages): + """Make destination_client.paginated_request(get)(...) yield ``pages`` in order.""" + inner = AsyncMock(side_effect=pages) + config.destination_client.paginated_request = MagicMock(return_value=inner) + return inner + + +def _post_paths(mock_config): + return [c.args[0] for c in mock_config.destination_client.post.call_args_list] + + +class TestV2CreatePath: + def test_flag_off_uses_v2_not_v1(self, mock_config): + """b: with the flag off, create goes through the v2 endpoint and never + touches the v1 API (preserves the default upstream behavior).""" + mock_config.use_v1_user_api = False + instance = Users(mock_config) + instance._existing_resources_map = {} + mock_config.destination_client.post = AsyncMock(return_value={"data": {"id": "dest-x", "attributes": {}}}) + + asyncio.run(instance.create_resource("src-a", _make_user("user-a@example.com", "shared@example.com", "src-a"))) + + assert _post_paths(mock_config) == ["/api/v2/users"] + + +class TestV1CreatePath: + def test_flag_on_uses_v1_with_handle_not_v2(self, mock_config): + """d1: with the flag on, create posts to /api/v1/user carrying the handle, + never calls the v2 users endpoint, and returns the reconciled v2 UUID.""" + mock_config.use_v1_user_api = True + instance = Users(mock_config) + instance._existing_resources_map = {} + dest_user = _make_user("user-a@example.com", "shared@example.com", "dest-uuid-a") + mock_config.destination_client.post = AsyncMock(return_value={"data": {}}) + _mock_paginated(mock_config, [[dest_user]]) + + _id, r = asyncio.run( + instance.create_resource("src-a", _make_user("user-a@example.com", "shared@example.com", "src-a")) + ) + + v1_calls = [c for c in mock_config.destination_client.post.call_args_list if c.args[0] == "/api/v1/user"] + assert len(v1_calls) == 1 + assert v1_calls[0].args[1]["handle"] == "user-a@example.com" + assert "/api/v2/users" not in _post_paths(mock_config) + assert r["id"] == "dest-uuid-a" + + def test_shared_email_distinct_handles_both_created_via_v1(self, mock_config): + """The core fix. Two users share an email but have distinct handles: one + handle equals the shared email, the other differs. Under v2 the first + create derives its handle from the email and steals the second user's + handle, so the second 409s and can never be created. Via v1 each user is + created with its OWN handle, so both succeed.""" + mock_config.use_v1_user_api = True + instance = Users(mock_config) + instance._existing_resources_map = {} + dest_a = _make_user("abc@example.com", "jsmith@example.com", "dest-a") + dest_b = _make_user("jsmith@example.com", "jsmith@example.com", "dest-b") + mock_config.destination_client.post = AsyncMock(return_value={"data": {}}) + _mock_paginated(mock_config, [[dest_a], [dest_b]]) + + asyncio.run(instance.create_resource("src-a", _make_user("abc@example.com", "jsmith@example.com", "src-a"))) + asyncio.run(instance.create_resource("src-b", _make_user("jsmith@example.com", "jsmith@example.com", "src-b"))) + + v1_handles = [ + c.args[1]["handle"] + for c in mock_config.destination_client.post.call_args_list + if c.args[0] == "/api/v1/user" + ] + assert v1_handles == ["abc@example.com", "jsmith@example.com"] + assert "/api/v2/users" not in _post_paths(mock_config) + + def test_requires_handle(self, mock_config): + """d0: v1 creation raises when there is no handle to send.""" + mock_config.use_v1_user_api = True + instance = Users(mock_config) + with pytest.raises(ValueError): + asyncio.run(instance._create_via_v1("", "Person Name", "e@example.com")) + + def test_v1_post_failure_reraises(self, mock_config): + """d4: a v1 POST failure propagates so the apply loop counts it failed + and continues (DR-safe).""" + mock_config.use_v1_user_api = True + instance = Users(mock_config) + instance._existing_resources_map = {} + mock_config.destination_client.post = AsyncMock(side_effect=_http_error(400)) + with pytest.raises(CustomClientHTTPError): + asyncio.run( + instance.create_resource("src-a", _make_user("user-a@example.com", "shared@example.com", "src-a")) + ) + + +class TestReconcileByHandle: + def test_no_exact_match_raises(self, mock_config): + """If no exact-handle match is ever found after the v1 create, it raises.""" + mock_config.use_v1_user_api = True + instance = Users(mock_config) + instance._existing_resources_map = {} + mock_config.destination_client.post = AsyncMock(return_value={"data": {}}) + other = _make_user("user-b@example.com", "shared@example.com", "dest-b") + _mock_paginated(mock_config, [[other], [other], [other]]) + with pytest.raises(ValueError): + asyncio.run( + instance.create_resource("src-a", _make_user("user-a@example.com", "shared@example.com", "src-a")) + ) + + def test_requeries_on_empty_then_matches(self, mock_config): + """d3b: read-after-write — an empty first page then a match re-queries + (exactly two calls) rather than giving up on the first empty result.""" + instance = Users(mock_config) + match = _make_user("user-a@example.com", "shared@example.com", "dest-uuid-a") + inner = _mock_paginated(mock_config, [[], [match]]) + user = asyncio.run(instance._get_destination_user_by_handle("user-a@example.com")) + assert user["id"] == "dest-uuid-a" + assert inner.call_count == 2 + + def test_selects_exact_case_handle_among_candidates(self, mock_config): + """d5: with multiple filter candidates, only the exact-case handle wins.""" + instance = Users(mock_config) + a = _make_user("user-a@example.com", "shared@example.com", "dest-a") + b = _make_user("user-b@example.com", "shared@example.com", "dest-b") + _mock_paginated(mock_config, [[b, a]]) + user = asyncio.run(instance._get_destination_user_by_handle("user-a@example.com")) + assert user["id"] == "dest-a" + + +class TestUpdatePathRegression: + def test_existing_handle_takes_update_path_no_create(self, mock_config): + """e: an existing destination handle routes to the update path — no v1 or + v2 create, no duplicate.""" + mock_config.use_v1_user_api = True + instance = Users(mock_config) + dest_user = _make_user("user-a@example.com", "shared@example.com", "dest-a") + instance._existing_resources_map = {"user-a@example.com": dest_user} + mock_config.destination_client.post = AsyncMock() + mock_config.destination_client.patch = AsyncMock(return_value={"data": dest_user}) + + _id, r = asyncio.run( + instance.create_resource("src-a", _make_user("user-a@example.com", "shared@example.com", "src-a")) + ) + mock_config.destination_client.post.assert_not_called() + assert r["id"] == "dest-a" From dd997305b3ea770b104a0a0c90cd7f0f045812f8 Mon Sep 17 00:00:00 2001 From: michael-richey <41595765+michael-richey@users.noreply.github.com> Date: Fri, 17 Jul 2026 15:16:18 -0400 Subject: [PATCH 2/7] fix(users): assign roles after v1 creation --- datadog_sync/model/users.py | 35 ++++++++++++-- tests/unit/test_users.py | 96 ++++++++++++++++++++++++++++++------- 2 files changed, 111 insertions(+), 20 deletions(-) diff --git a/datadog_sync/model/users.py b/datadog_sync/model/users.py index 52403428e..266280569 100644 --- a/datadog_sync/model/users.py +++ b/datadog_sync/model/users.py @@ -94,7 +94,10 @@ async def create_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: # collapses distinct-handle users that share an email onto one handle. # v1 accepts an explicit handle, so each user keeps its own. return _id, await self._create_via_v1( - attributes.get("handle"), attributes.get("name"), attributes.get("email") + attributes.get("handle"), + attributes.get("name"), + attributes.get("email"), + resource.get("relationships", {}).get("roles", {}).get("data", []), ) destination_client = self.config.destination_client @@ -104,7 +107,13 @@ async def create_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: resp = await destination_client.post(self.resource_config.base_path, {"data": resource}) return _id, resp["data"] - async def _create_via_v1(self, handle: Optional[str], name: Optional[str], email: Optional[str]) -> Dict: + async def _create_via_v1( + self, + handle: Optional[str], + name: Optional[str], + email: Optional[str], + desired_roles: Optional[List[Dict]] = None, + ) -> Dict: """Create the user via the v1 API, which accepts an explicit handle. v2 ``POST /api/v2/users`` cannot set a handle (it is derived from the @@ -124,8 +133,26 @@ async def _create_via_v1(self, handle: Optional[str], name: Optional[str], email user = await self._get_destination_user_by_handle(handle) if user is None: raise ValueError("v1-created user not found by handle after create") + await self._assign_missing_roles(user, desired_roles or []) return user + async def _assign_missing_roles(self, user: Dict, desired_roles: List[Dict]) -> None: + """Assign missing roles and keep the reconciled user state accurate.""" + existing_roles = user.setdefault("relationships", {}).setdefault("roles", {}).setdefault("data", []) + existing_role_ids = { + role["id"] for role in existing_roles if isinstance(role, dict) and role.get("id") is not None + } + + for role in desired_roles: + if not isinstance(role, dict) or role.get("id") is None: + continue + role_id = role["id"] + if role_id in existing_role_ids: + continue + if await self.add_user_to_role(user["id"], role_id): + existing_roles.append(dict(role)) + existing_role_ids.add(role_id) + async def _get_destination_user_by_handle(self, handle: str) -> Optional[Dict]: """Return the destination user whose handle matches exactly, or None. @@ -189,13 +216,15 @@ async def update_user_roles(self, _id, diff): role_id = new_val["id"] if isinstance(new_val, dict) else new_val await self.add_user_to_role(_id, role_id) - async def add_user_to_role(self, user_id, role_id): + async def add_user_to_role(self, user_id, role_id) -> bool: destination_client = self.config.destination_client payload = {"data": {"id": user_id, "type": "users"}} try: await destination_client.post(self.roles_path.format(role_id), payload) + return True except CustomClientHTTPError as e: self.config.logger.error("error adding user: %s to role %s: %s", user_id, role_id, e) + return False async def remove_user_from_role(self, user_id, role_id): destination_client = self.config.destination_client diff --git a/tests/unit/test_users.py b/tests/unit/test_users.py index 8a74a1db8..cc1df7681 100644 --- a/tests/unit/test_users.py +++ b/tests/unit/test_users.py @@ -39,7 +39,7 @@ def _http_error(status): return CustomClientHTTPError(response) -def _make_user(handle, email, user_id, name="User"): +def _make_user(handle, email, user_id, name="User", roles=None): """Build a user dict shaped like the v2 GET/create response.""" return { "id": user_id, @@ -50,7 +50,7 @@ def _make_user(handle, email, user_id, name="User"): "name": name, "disabled": False, }, - "relationships": {"roles": {"data": []}}, + "relationships": {"roles": {"data": roles or []}}, } @@ -123,9 +123,7 @@ def test_v2_post_body_excludes_handle_and_disabled(self, mock_config): mock_config.use_v1_user_api = False instance = Users(mock_config) instance._existing_resources_map = {} - mock_config.destination_client.post = AsyncMock( - return_value={"data": {"id": "dest-x", "attributes": {}}} - ) + mock_config.destination_client.post = AsyncMock(return_value={"data": {"id": "dest-x", "attributes": {}}}) src = _make_user("user-a@example.com", "shared@example.com", "src-a") asyncio.run(instance.create_resource("src-a", src)) @@ -140,14 +138,10 @@ def test_v2_patch_body_excludes_handle(self, mock_config): """f2: the v2 update (PATCH) body must not carry the read-only handle.""" instance = Users(mock_config) _id = "src-a" - mock_config.state.destination["users"][_id] = _make_user( - "user-a@example.com", "shared@example.com", "dest-a" - ) + mock_config.state.destination["users"][_id] = _make_user("user-a@example.com", "shared@example.com", "dest-a") # A differing name forces a diff -> the PATCH branch. src = _make_user("user-a@example.com", "shared@example.com", "dest-a", name="New Name") - mock_config.destination_client.patch = AsyncMock( - return_value={"data": {"id": "dest-a", "attributes": {}}} - ) + mock_config.destination_client.patch = AsyncMock(return_value={"data": {"id": "dest-a", "attributes": {}}}) asyncio.run(instance.update_resource(_id, src)) @@ -178,16 +172,12 @@ def test_build_config_flag_default_false(self, tmp_path): def test_sync_accepts_v1_flag(self): """p2: the sync command recognizes the flag (exit 2 == usage error).""" - result = CliRunner(mix_stderr=False).invoke( - cli, ["sync", "--use-v1-user-api=true", "--validate=false"] - ) + result = CliRunner(mix_stderr=False).invoke(cli, ["sync", "--use-v1-user-api=true", "--validate=false"]) assert result.exit_code != 2 def test_migrate_accepts_v1_flag(self): """p3: migrate reuses @sync_options, so it recognizes the flag too.""" - result = CliRunner(mix_stderr=False).invoke( - cli, ["migrate", "--use-v1-user-api=true", "--validate=false"] - ) + result = CliRunner(mix_stderr=False).invoke(cli, ["migrate", "--use-v1-user-api=true", "--validate=false"]) assert result.exit_code != 2 @@ -262,6 +252,78 @@ def test_shared_email_distinct_handles_both_created_via_v1(self, mock_config): assert v1_handles == ["abc@example.com", "jsmith@example.com"] assert "/api/v2/users" not in _post_paths(mock_config) + def test_assigns_only_missing_roles_and_returns_updated_state(self, mock_config): + """v1 create assigns mapped roles after resolving the v2 UUID. + + Roles already present on the reconciled user are not posted again, and + successful assignments are reflected in the destination state returned + by create_resource. + """ + mock_config.use_v1_user_api = True + instance = Users(mock_config) + instance._existing_resources_map = {} + existing_role = {"id": "role-dst-existing", "type": "roles"} + missing_role = {"id": "role-dst-missing", "type": "roles"} + dest_user = _make_user( + "user-a@example.com", + "shared@example.com", + "dest-uuid-a", + roles=[existing_role], + ) + source_user = _make_user( + "user-a@example.com", + "shared@example.com", + "src-a", + roles=[existing_role, missing_role], + ) + mock_config.destination_client.post = AsyncMock(return_value={"data": {}}) + _mock_paginated(mock_config, [[dest_user]]) + + _, created = asyncio.run(instance.create_resource("src-a", source_user)) + + assert _post_paths(mock_config) == [ + "/api/v1/user", + "/api/v2/roles/role-dst-missing/users", + ] + assert created["relationships"]["roles"]["data"] == [existing_role, missing_role] + + def test_role_failure_does_not_block_remaining_roles_or_corrupt_state(self, mock_config): + """A failed role assignment is logged while later roles are attempted. + + Only successful assignments are recorded in returned destination state, + leaving failed roles eligible for retry on a later sync. + """ + mock_config.use_v1_user_api = True + instance = Users(mock_config) + instance._existing_resources_map = {} + failed_role = {"id": "role-dst-failed", "type": "roles"} + successful_role = {"id": "role-dst-success", "type": "roles"} + dest_user = _make_user("user-a@example.com", "shared@example.com", "dest-uuid-a") + source_user = _make_user( + "user-a@example.com", + "shared@example.com", + "src-a", + roles=[failed_role, successful_role], + ) + + async def post(path, _body): + if path == "/api/v2/roles/role-dst-failed/users": + raise _http_error(403) + return {"data": {}} + + mock_config.destination_client.post = AsyncMock(side_effect=post) + _mock_paginated(mock_config, [[dest_user]]) + + _, created = asyncio.run(instance.create_resource("src-a", source_user)) + + assert _post_paths(mock_config) == [ + "/api/v1/user", + "/api/v2/roles/role-dst-failed/users", + "/api/v2/roles/role-dst-success/users", + ] + assert created["relationships"]["roles"]["data"] == [successful_role] + mock_config.logger.error.assert_called_once() + def test_requires_handle(self, mock_config): """d0: v1 creation raises when there is no handle to send.""" mock_config.use_v1_user_api = True From 4073c41a49cb232105ad029676a4dbb500035694 Mon Sep 17 00:00:00 2001 From: michael-richey <41595765+michael-richey@users.noreply.github.com> Date: Fri, 17 Jul 2026 15:31:34 -0400 Subject: [PATCH 3/7] fix(users): report partial role assignment failures --- datadog_sync/model/users.py | 37 +++++++++++++++++++++++++++++++------ tests/unit/test_users.py | 16 +++++++++++----- 2 files changed, 42 insertions(+), 11 deletions(-) diff --git a/datadog_sync/model/users.py b/datadog_sync/model/users.py index 266280569..7821dd95e 100644 --- a/datadog_sync/model/users.py +++ b/datadog_sync/model/users.py @@ -14,6 +14,17 @@ from datadog_sync.utils.custom_client import CustomClient +class UserRoleAssignmentError(RuntimeError): + """Role assignment failed after the user was created and reconciled.""" + + def __init__(self, user: Dict, failed_role_ids: List[str]) -> None: + self.user = user + self.failed_role_ids = tuple(failed_role_ids) + count = len(failed_role_ids) + assignment = "role assignment" if count == 1 else "role assignments" + super().__init__(f"{count} {assignment} failed after v1 user creation") + + class Users(BaseResource): resource_type = "users" resource_config = ResourceConfig( @@ -93,12 +104,20 @@ async def create_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: # v2 create cannot set a handle (it derives one from the email), which # collapses distinct-handle users that share an email onto one handle. # v1 accepts an explicit handle, so each user keeps its own. - return _id, await self._create_via_v1( - attributes.get("handle"), - attributes.get("name"), - attributes.get("email"), - resource.get("relationships", {}).get("roles", {}).get("data", []), - ) + try: + user = await self._create_via_v1( + attributes.get("handle"), + attributes.get("name"), + attributes.get("email"), + resource.get("relationships", {}).get("roles", {}).get("data", []), + ) + except UserRoleAssignmentError as e: + # Preserve the reconciled UUID and successful memberships so + # downstream resources can still resolve this user. Re-raising + # makes the apply handler count and emit the partial failure. + self.config.state.destination[self.resource_type][_id] = e.user + raise + return _id, user destination_client = self.config.destination_client # handle is read-only in v2 (derived from email) and must not be sent. @@ -142,6 +161,7 @@ async def _assign_missing_roles(self, user: Dict, desired_roles: List[Dict]) -> existing_role_ids = { role["id"] for role in existing_roles if isinstance(role, dict) and role.get("id") is not None } + failed_role_ids = [] for role in desired_roles: if not isinstance(role, dict) or role.get("id") is None: @@ -152,6 +172,11 @@ async def _assign_missing_roles(self, user: Dict, desired_roles: List[Dict]) -> if await self.add_user_to_role(user["id"], role_id): existing_roles.append(dict(role)) existing_role_ids.add(role_id) + else: + failed_role_ids.append(str(role_id)) + + if failed_role_ids: + raise UserRoleAssignmentError(user, failed_role_ids) async def _get_destination_user_by_handle(self, handle: str) -> Optional[Dict]: """Return the destination user whose handle matches exactly, or None. diff --git a/tests/unit/test_users.py b/tests/unit/test_users.py index cc1df7681..afe203669 100644 --- a/tests/unit/test_users.py +++ b/tests/unit/test_users.py @@ -26,7 +26,7 @@ from datadog_sync.cli import cli from datadog_sync.constants import Command -from datadog_sync.model.users import Users +from datadog_sync.model.users import UserRoleAssignmentError, Users from datadog_sync.utils.configuration import build_config from datadog_sync.utils.resource_utils import CustomClientHTTPError, check_diff @@ -287,11 +287,12 @@ def test_assigns_only_missing_roles_and_returns_updated_state(self, mock_config) ] assert created["relationships"]["roles"]["data"] == [existing_role, missing_role] - def test_role_failure_does_not_block_remaining_roles_or_corrupt_state(self, mock_config): - """A failed role assignment is logged while later roles are attempted. + def test_role_failure_persists_partial_state_and_reports_failure(self, mock_config): + """A failed role assignment is reported after later roles are attempted. Only successful assignments are recorded in returned destination state, - leaving failed roles eligible for retry on a later sync. + leaving failed roles eligible for retry on a later sync. The exception + lets the apply handler count the otherwise-partial create as a failure. """ mock_config.use_v1_user_api = True instance = Users(mock_config) @@ -314,13 +315,18 @@ async def post(path, _body): mock_config.destination_client.post = AsyncMock(side_effect=post) _mock_paginated(mock_config, [[dest_user]]) - _, created = asyncio.run(instance.create_resource("src-a", source_user)) + with pytest.raises( + UserRoleAssignmentError, match="1 role assignment failed after v1 user creation" + ) as exc_info: + asyncio.run(instance.create_resource("src-a", source_user)) + assert exc_info.value.failed_role_ids == ("role-dst-failed",) assert _post_paths(mock_config) == [ "/api/v1/user", "/api/v2/roles/role-dst-failed/users", "/api/v2/roles/role-dst-success/users", ] + created = mock_config.state.destination["users"]["src-a"] assert created["relationships"]["roles"]["data"] == [successful_role] mock_config.logger.error.assert_called_once() From a523afeb2d4194fa4eee5d8868088dd440d03ce7 Mon Sep 17 00:00:00 2001 From: michael-richey <41595765+michael-richey@users.noreply.github.com> Date: Fri, 17 Jul 2026 15:58:31 -0400 Subject: [PATCH 4/7] fix(users): report retried role assignment failures --- datadog_sync/model/users.py | 59 +++++++++++++++++++++---------------- tests/unit/test_users.py | 42 +++++++++++++++++++++++++- 2 files changed, 75 insertions(+), 26 deletions(-) diff --git a/datadog_sync/model/users.py b/datadog_sync/model/users.py index 7821dd95e..a8833fa08 100644 --- a/datadog_sync/model/users.py +++ b/datadog_sync/model/users.py @@ -15,14 +15,14 @@ class UserRoleAssignmentError(RuntimeError): - """Role assignment failed after the user was created and reconciled.""" + """One or more role assignments failed while reconciling a user.""" def __init__(self, user: Dict, failed_role_ids: List[str]) -> None: self.user = user self.failed_role_ids = tuple(failed_role_ids) count = len(failed_role_ids) assignment = "role assignment" if count == 1 else "role assignments" - super().__init__(f"{count} {assignment} failed after v1 user creation") + super().__init__(f"{count} {assignment} failed while reconciling user") class Users(BaseResource): @@ -200,19 +200,46 @@ async def _get_destination_user_by_handle(self, handle: str) -> Optional[Dict]: async def update_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: destination_client = self.config.destination_client - diff = check_diff(self.resource_config, self.config.state.destination[self.resource_type][_id], resource) + destination_user = self.config.state.destination[self.resource_type][_id] + diff = check_diff(self.resource_config, destination_user, resource) if diff: - await self.update_user_roles(self.config.state.destination[self.resource_type][_id]["id"], diff) - resource["id"] = self.config.state.destination[self.resource_type][_id]["id"] + role_error = None + try: + await self._assign_missing_roles( + destination_user, + resource.get("relationships", {}).get("roles", {}).get("data", []), + ) + except UserRoleAssignmentError as e: + # Continue with unrelated attribute updates, then report the + # partial role failure so the apply handler does not emit success. + role_error = e + + resource["id"] = destination_user["id"] resource.pop("relationships", None) resource["attributes"].pop("handle", None) resp = await destination_client.patch( - self.resource_config.base_path + f"/{self.config.state.destination[self.resource_type][_id]['id']}", + self.resource_config.base_path + f"/{destination_user['id']}", {"data": resource}, ) + if role_error is not None: + updated_user = resp["data"] + updated_roles = ( + updated_user.setdefault("relationships", {}).setdefault("roles", {}).setdefault("data", []) + ) + updated_role_ids = { + role["id"] for role in updated_roles if isinstance(role, dict) and role.get("id") is not None + } + for role in destination_user.get("relationships", {}).get("roles", {}).get("data", []): + role_id = role.get("id") if isinstance(role, dict) else None + if role_id is not None and role_id not in updated_role_ids: + updated_roles.append(dict(role)) + updated_role_ids.add(role_id) + self.config.state.destination[self.resource_type][_id] = updated_user + role_error.user = updated_user + raise role_error return _id, resp["data"] - return _id, self.config.state.destination[self.resource_type][_id] + return _id, destination_user async def delete_resource(self, _id: str) -> None: destination_client = self.config.destination_client @@ -223,24 +250,6 @@ async def delete_resource(self, _id: str) -> None: def connect_id(self, key: str, r_obj: Dict, resource_to_connect: str) -> Optional[List[str]]: return super(Users, self).connect_id(key, r_obj, resource_to_connect) - async def update_user_roles(self, _id, diff): - for k, v in diff.items(): - if k == "iterable_item_added": - for key, value in diff["iterable_item_added"].items(): - if "roles" in key: - await self.add_user_to_role(_id, value["id"]) - # elif k == "iterable_item_removed": - # for key, value in diff["iterable_item_removed"].items(): - # if "roles" in key: - # await self.remove_user_from_role(_id, value["id"]) - elif k == "values_changed": - for key, value in diff["values_changed"].items(): - if "roles" in key: - # await self.remove_user_from_role(_id, value["old_value"]) - new_val = value["new_value"] - role_id = new_val["id"] if isinstance(new_val, dict) else new_val - await self.add_user_to_role(_id, role_id) - async def add_user_to_role(self, user_id, role_id) -> bool: destination_client = self.config.destination_client payload = {"data": {"id": user_id, "type": "users"}} diff --git a/tests/unit/test_users.py b/tests/unit/test_users.py index afe203669..8e850b4f0 100644 --- a/tests/unit/test_users.py +++ b/tests/unit/test_users.py @@ -316,7 +316,7 @@ async def post(path, _body): _mock_paginated(mock_config, [[dest_user]]) with pytest.raises( - UserRoleAssignmentError, match="1 role assignment failed after v1 user creation" + UserRoleAssignmentError, match="1 role assignment failed while reconciling user" ) as exc_info: asyncio.run(instance.create_resource("src-a", source_user)) @@ -385,6 +385,46 @@ def test_selects_exact_case_handle_among_candidates(self, mock_config): class TestUpdatePathRegression: + def test_role_retry_persists_partial_state_and_reports_failure(self, mock_config): + """A later run retries missing roles without reporting full success.""" + instance = Users(mock_config) + failed_role = {"id": "role-dst-failed", "type": "roles"} + successful_role = {"id": "role-dst-success", "type": "roles"} + dest_user = _make_user("user-a@example.com", "shared@example.com", "dest-a") + source_user = _make_user( + "user-a@example.com", + "shared@example.com", + "src-a", + name="Updated User", + roles=[failed_role, successful_role], + ) + mock_config.state.destination["users"]["src-a"] = dest_user + instance.add_user_to_role = AsyncMock(side_effect=[False, True]) + updated_user = { + "id": "dest-a", + "type": "users", + "attributes": { + "handle": "user-a@example.com", + "email": "shared@example.com", + "name": "Updated User", + "disabled": False, + }, + } + mock_config.destination_client.patch = AsyncMock(return_value={"data": updated_user}) + + with pytest.raises(UserRoleAssignmentError) as exc_info: + asyncio.run(instance.update_resource("src-a", source_user)) + + assert exc_info.value.failed_role_ids == ("role-dst-failed",) + assert [call.args for call in instance.add_user_to_role.await_args_list] == [ + ("dest-a", "role-dst-failed"), + ("dest-a", "role-dst-success"), + ] + mock_config.destination_client.patch.assert_awaited_once() + stored = mock_config.state.destination["users"]["src-a"] + assert stored["attributes"]["name"] == "Updated User" + assert stored["relationships"]["roles"]["data"] == [successful_role] + def test_existing_handle_takes_update_path_no_create(self, mock_config): """e: an existing destination handle routes to the update path — no v1 or v2 create, no duplicate.""" From 9d48c849d55c6bf646f656218c2afd7452c67638 Mon Sep 17 00:00:00 2001 From: michael-richey <41595765+michael-richey@users.noreply.github.com> Date: Fri, 17 Jul 2026 16:00:19 -0400 Subject: [PATCH 5/7] chore(tests): remove unused coverage dependency --- setup.cfg | 1 - 1 file changed, 1 deletion(-) diff --git a/setup.cfg b/setup.cfg index e9bc0462e..0b06fe831 100644 --- a/setup.cfg +++ b/setup.cfg @@ -56,7 +56,6 @@ tests = black==24.3.0 pytest==8.1.1 pytest-black - pytest-cov pytest-console-scripts pytest-recording==0.13.2 pytest-retry==1.7.0 From 8b3fb76d44a1c5ad2fdb5aac6b9d56029c534302 Mon Sep 17 00:00:00 2001 From: michael-richey <41595765+michael-richey@users.noreply.github.com> Date: Fri, 17 Jul 2026 16:07:20 -0400 Subject: [PATCH 6/7] fix(users): persist successful role retries --- datadog_sync/model/users.py | 29 ++++++++++++++++------------- tests/unit/test_users.py | 36 ++++++++++++++++++++++++++++++++++++ 2 files changed, 52 insertions(+), 13 deletions(-) diff --git a/datadog_sync/model/users.py b/datadog_sync/model/users.py index a8833fa08..6b9a947ff 100644 --- a/datadog_sync/model/users.py +++ b/datadog_sync/model/users.py @@ -178,6 +178,20 @@ async def _assign_missing_roles(self, user: Dict, desired_roles: List[Dict]) -> if failed_role_ids: raise UserRoleAssignmentError(user, failed_role_ids) + @staticmethod + def _merge_role_state(updated_user: Dict, role_state: Dict) -> Dict: + """Preserve known role memberships when a user PATCH omits them.""" + updated_roles = updated_user.setdefault("relationships", {}).setdefault("roles", {}).setdefault("data", []) + updated_role_ids = { + role["id"] for role in updated_roles if isinstance(role, dict) and role.get("id") is not None + } + for role in role_state.get("relationships", {}).get("roles", {}).get("data", []): + role_id = role.get("id") if isinstance(role, dict) else None + if role_id is not None and role_id not in updated_role_ids: + updated_roles.append(dict(role)) + updated_role_ids.add(role_id) + return updated_user + async def _get_destination_user_by_handle(self, handle: str) -> Optional[Dict]: """Return the destination user whose handle matches exactly, or None. @@ -221,24 +235,13 @@ async def update_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: self.resource_config.base_path + f"/{destination_user['id']}", {"data": resource}, ) + updated_user = self._merge_role_state(resp["data"], destination_user) if role_error is not None: - updated_user = resp["data"] - updated_roles = ( - updated_user.setdefault("relationships", {}).setdefault("roles", {}).setdefault("data", []) - ) - updated_role_ids = { - role["id"] for role in updated_roles if isinstance(role, dict) and role.get("id") is not None - } - for role in destination_user.get("relationships", {}).get("roles", {}).get("data", []): - role_id = role.get("id") if isinstance(role, dict) else None - if role_id is not None and role_id not in updated_role_ids: - updated_roles.append(dict(role)) - updated_role_ids.add(role_id) self.config.state.destination[self.resource_type][_id] = updated_user role_error.user = updated_user raise role_error - return _id, resp["data"] + return _id, updated_user return _id, destination_user async def delete_resource(self, _id: str) -> None: diff --git a/tests/unit/test_users.py b/tests/unit/test_users.py index 8e850b4f0..b688d9ba7 100644 --- a/tests/unit/test_users.py +++ b/tests/unit/test_users.py @@ -385,6 +385,42 @@ def test_selects_exact_case_handle_among_candidates(self, mock_config): class TestUpdatePathRegression: + def test_successful_role_retry_remains_in_persisted_state(self, mock_config): + """A successful role retry is not lost when PATCH omits relationships.""" + instance = Users(mock_config) + successful_role = {"id": "role-dst-success", "type": "roles"} + dest_user = _make_user("user-a@example.com", "shared@example.com", "dest-a") + mock_config.state.destination["users"]["src-a"] = dest_user + instance.add_user_to_role = AsyncMock(return_value=True) + updated_user = { + "id": "dest-a", + "type": "users", + "attributes": { + "handle": "user-a@example.com", + "email": "shared@example.com", + "name": "Updated User", + "disabled": False, + }, + } + mock_config.destination_client.patch = AsyncMock(return_value={"data": updated_user}) + + def source_user(): + return _make_user( + "user-a@example.com", + "shared@example.com", + "src-a", + name="Updated User", + roles=[successful_role], + ) + + asyncio.run(instance._update_resource("src-a", source_user())) + asyncio.run(instance._update_resource("src-a", source_user())) + + stored = mock_config.state.destination["users"]["src-a"] + assert stored["relationships"]["roles"]["data"] == [successful_role] + instance.add_user_to_role.assert_awaited_once_with("dest-a", "role-dst-success") + mock_config.destination_client.patch.assert_awaited_once() + def test_role_retry_persists_partial_state_and_reports_failure(self, mock_config): """A later run retries missing roles without reporting full success.""" instance = Users(mock_config) From 60ed270cc491e32108143f6affeee3ef9ab66915 Mon Sep 17 00:00:00 2001 From: michael-richey <41595765+michael-richey@users.noreply.github.com> Date: Fri, 17 Jul 2026 16:10:02 -0400 Subject: [PATCH 7/7] fix(users): delay reconciliation retries --- datadog_sync/model/users.py | 10 +++++++--- tests/unit/test_users.py | 18 +++++++++++------- 2 files changed, 18 insertions(+), 10 deletions(-) diff --git a/datadog_sync/model/users.py b/datadog_sync/model/users.py index 6b9a947ff..be18e046a 100644 --- a/datadog_sync/model/users.py +++ b/datadog_sync/model/users.py @@ -4,6 +4,7 @@ # Copyright 2019 Datadog, Inc. from __future__ import annotations +import asyncio from typing import TYPE_CHECKING, Any, Optional, List, Dict, Tuple, cast from datadog_sync.utils.base_resource import BaseResource, ResourceConfig @@ -66,6 +67,7 @@ class Users(BaseResource): page_size=500, ) roles_path: str = "/api/v2/roles/{}/users" + user_lookup_retry_delays: Tuple[float, ...] = (1.0, 2.0) async def get_resources(self, client: CustomClient) -> List[Dict]: resp = await client.paginated_request(client.get)( @@ -196,11 +198,11 @@ async def _get_destination_user_by_handle(self, handle: str) -> Optional[Dict]: """Return the destination user whose handle matches exactly, or None. Transient HTTP errors are already retried by the client's - ``request_with_retry``; this adds a small, sleep-free re-query loop only - to absorb read-after-write visibility lag after a v1 create. + ``request_with_retry``; this adds a bounded re-query loop with delays to + absorb read-after-write visibility lag after a v1 create. """ destination_client = self.config.destination_client - for _ in range(3): + for attempt in range(len(self.user_lookup_retry_delays) + 1): resp = await destination_client.paginated_request(destination_client.get)( self.resource_config.base_path, pagination_config=self.pagination_config, @@ -209,6 +211,8 @@ async def _get_destination_user_by_handle(self, handle: str) -> Optional[Dict]: for user in resp: if user.get("attributes", {}).get("handle") == handle: return user + if attempt < len(self.user_lookup_retry_delays): + await asyncio.sleep(self.user_lookup_retry_delays[attempt]) return None async def update_resource(self, _id: str, resource: Dict) -> Tuple[str, Dict]: diff --git a/tests/unit/test_users.py b/tests/unit/test_users.py index b688d9ba7..74417b996 100644 --- a/tests/unit/test_users.py +++ b/tests/unit/test_users.py @@ -19,7 +19,7 @@ """ import asyncio -from unittest.mock import AsyncMock, MagicMock +from unittest.mock import AsyncMock, MagicMock, call, patch import pytest from click.testing import CliRunner @@ -359,20 +359,24 @@ def test_no_exact_match_raises(self, mock_config): mock_config.destination_client.post = AsyncMock(return_value={"data": {}}) other = _make_user("user-b@example.com", "shared@example.com", "dest-b") _mock_paginated(mock_config, [[other], [other], [other]]) - with pytest.raises(ValueError): - asyncio.run( - instance.create_resource("src-a", _make_user("user-a@example.com", "shared@example.com", "src-a")) - ) + with patch("datadog_sync.model.users.asyncio.sleep", new_callable=AsyncMock) as sleep: + with pytest.raises(ValueError): + asyncio.run( + instance.create_resource("src-a", _make_user("user-a@example.com", "shared@example.com", "src-a")) + ) + assert sleep.await_args_list == [call(1.0), call(2.0)] def test_requeries_on_empty_then_matches(self, mock_config): """d3b: read-after-write — an empty first page then a match re-queries - (exactly two calls) rather than giving up on the first empty result.""" + after a bounded delay rather than immediately retrying.""" instance = Users(mock_config) match = _make_user("user-a@example.com", "shared@example.com", "dest-uuid-a") inner = _mock_paginated(mock_config, [[], [match]]) - user = asyncio.run(instance._get_destination_user_by_handle("user-a@example.com")) + with patch("datadog_sync.model.users.asyncio.sleep", new_callable=AsyncMock) as sleep: + user = asyncio.run(instance._get_destination_user_by_handle("user-a@example.com")) assert user["id"] == "dest-uuid-a" assert inner.call_count == 2 + assert sleep.await_args_list == [call(1.0)] def test_selects_exact_case_handle_among_candidates(self, mock_config): """d5: with multiple filter candidates, only the exact-case handle wins."""