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

1from datetime import UTC, datetime, timedelta 

2from logging import getLogger 

3import re 

4 

5import httpx 

6 

7from atlas_ecommerce.clients.delivery import PermanentDeliveryError, TransientDeliveryError, raise_for_response 

8from atlas_ecommerce.purchases import PurchaseEvent 

9 

10logger = getLogger(__name__) 

11 

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"} 

18 

19 

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) 

23 

24 

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]) 

28 

29 

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") 

34 

35 

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 

65 

66 

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] 

73 

74 

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) 

84 

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 ) 

112 

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 

118 

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 

140 

141 async def close(self) -> None: 

142 await self._http.aclose()