Coverage for atlas_ecommerce/clients/google_ads.py: 94%
78 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-25 16:36 +0000
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-25 16:36 +0000
1from datetime import UTC, datetime, timedelta
2from logging import getLogger
3import re
5import httpx
7from atlas_ecommerce.clients.delivery import PermanentDeliveryError, TransientDeliveryError, raise_for_response
8from atlas_ecommerce.purchases import PurchaseEvent
10logger = getLogger(__name__)
12TOKEN_URL = "https://oauth2.googleapis.com/token"
13INGEST_URL = "https://datamanager.googleapis.com/v1/events:ingest"
14CLICK_ID_KEYS = ("gclid", "gbraid", "wbraid")
15# Uploads are accepted up to the longest click-through window a conversion action can have.
16MAX_EVENT_AGE = timedelta(days=90)
17CONSENT_GRANTED = {"adUserData": "CONSENT_GRANTED", "adPersonalization": "CONSENT_GRANTED"}
20def digits_only(value: str) -> str:
21 """The UI shows 123-456-7890; the API wants the digits only."""
22 return re.sub(r"\D", "", value)
25def conversion_action_id(value: str) -> str:
26 """The bare id, or the trailing id of a customers/<cid>/conversionActions/<id> resource name."""
27 return digits_only(value.rsplit("/", 1)[-1])
30def to_rfc3339(moment: datetime) -> str:
31 if moment.tzinfo is None:
32 moment = moment.replace(tzinfo=UTC)
33 return moment.astimezone(UTC).strftime("%Y-%m-%dT%H:%M:%SZ")
36def ingest_request(
37 purchase: PurchaseEvent,
38 *,
39 customer_id: str,
40 conversion_action: str,
41 login_customer_id: str | None = None,
42 validate_only: bool = False,
43) -> dict:
44 destination: dict = {
45 "operatingAccount": {"accountType": "GOOGLE_ADS", "accountId": digits_only(customer_id)},
46 "productDestinationId": conversion_action_id(conversion_action),
47 }
48 if login_customer_id:
49 destination["loginAccount"] = {"accountType": "GOOGLE_ADS", "accountId": digits_only(login_customer_id)}
50 click_ids = purchase.attribution.click_ids
51 ad_identifiers = next(({key: click_ids[key]} for key in CLICK_ID_KEYS if click_ids.get(key)), {})
52 event = {
53 "adIdentifiers": ad_identifiers,
54 "eventTimestamp": to_rfc3339(purchase.happened_at),
55 "transactionId": purchase.order_id,
56 "conversionValue": float(purchase.value) if purchase.value else 0.0,
57 "currency": purchase.currency,
58 "eventSource": "WEB",
59 "consent": CONSENT_GRANTED,
60 }
61 body: dict = {"destinations": [destination], "events": [event]}
62 if validate_only:
63 body["validateOnly"] = True
64 return body
67def error_message(response: httpx.Response) -> str:
68 try:
69 error = response.json().get("error") or {}
70 except ValueError:
71 return response.text[:300]
72 return error.get("message") or response.text[:300]
75class GoogleDataManagerClient:
76 def __init__(
77 self, *, client_id: str, client_secret: str, refresh_token: str, http: httpx.AsyncClient | None = None
78 ) -> None:
79 self._client_id = client_id
80 self._client_secret = client_secret
81 self._refresh_token = refresh_token
82 self._access_token: str | None = None
83 self._http = http or httpx.AsyncClient(timeout=20.0)
85 async def send(self, payload: dict) -> None:
86 response = await self._post(payload, await self._token())
87 if response.status_code == 401:
88 response = await self._post(payload, await self._token(force_refresh=True))
89 if response.status_code in (401, 403):
90 # Credentials or account access to repair by an operator; the event itself is fine.
91 raise TransientDeliveryError(
92 f"Data Manager refused the credentials ({response.status_code}): {error_message(response)}"
93 )
94 if response.is_client_error and response.status_code not in (408, 429):
95 raise PermanentDeliveryError(
96 f"Data Manager rejected the request ({response.status_code}): {error_message(response)}"
97 )
98 raise_for_response(response)
99 body = response.json() if response.content else {}
100 for warning in body.get("fieldWarnings") or []:
101 logger.warning(
102 "Data Manager warning on %s: %s (%s)",
103 warning.get("field"),
104 warning.get("description"),
105 warning.get("reason"),
106 )
107 logger.info("Data Manager accepted the event; request %s", body.get("requestId"))
108 if payload.get("validateOnly"):
109 raise TransientDeliveryError(
110 "Data Manager validated the event without recording it; delivery stays pending"
111 )
113 async def _post(self, payload: dict, token: str) -> httpx.Response:
114 try:
115 return await self._http.post(INGEST_URL, json=payload, headers={"Authorization": f"Bearer {token}"})
116 except httpx.TransportError as exc:
117 raise TransientDeliveryError(f"Data Manager API unreachable: {exc}") from exc
119 async def _token(self, *, force_refresh: bool = False) -> str:
120 if self._access_token and not force_refresh:
121 return self._access_token
122 try:
123 response = await self._http.post(
124 TOKEN_URL,
125 data={
126 "grant_type": "refresh_token",
127 "refresh_token": self._refresh_token,
128 "client_id": self._client_id,
129 "client_secret": self._client_secret,
130 },
131 )
132 except httpx.TransportError as exc:
133 raise TransientDeliveryError(f"Google OAuth endpoint unreachable: {exc}") from exc
134 if not response.is_success:
135 # invalid_grant arrives as 400: the refresh token was revoked, and only re-authorising fixes that.
136 raise TransientDeliveryError(f"Google token refresh failed: {response.status_code} {response.text[:300]}")
137 self._access_token = response.json()["access_token"]
138 logger.info("Google access token refreshed")
139 return self._access_token
141 async def close(self) -> None:
142 await self._http.aclose()