Coverage for atlas_ecommerce/clients/linkedin_conversions.py: 92%

48 statements  

« prev     ^ index     » next       coverage.py v7.16.0, created at 2026-09-25 16:36 +0000

1from datetime import timedelta 

2from logging import getLogger 

3 

4import httpx 

5 

6from atlas_ecommerce.clients.delivery import TransientDeliveryError, raise_for_response 

7from atlas_ecommerce.purchases import PurchaseEvent 

8 

9logger = getLogger(__name__) 

10 

11CONVERSION_EVENTS_URL = "https://api.linkedin.com/rest/conversionEvents" 

12TOKEN_URL = "https://www.linkedin.com/oauth/v2/accessToken" 

13CLICK_ID_TYPE = "LINKEDIN_FIRST_PARTY_ADS_TRACKING_UUID" 

14MAX_EVENT_AGE = timedelta(days=90) 

15 

16 

17def purchase_event(purchase: PurchaseEvent, conversion_urn: str) -> dict: 

18 return { 

19 "conversion": conversion_urn, 

20 "conversionHappenedAt": int(purchase.happened_at.timestamp() * 1000), 

21 "conversionValue": {"currencyCode": purchase.currency, "amount": purchase.value}, 

22 "user": {"userIds": [{"idType": CLICK_ID_TYPE, "idValue": purchase.attribution.click_ids["li_fat_id"]}]}, 

23 "eventId": purchase.order_id, 

24 } 

25 

26 

27class LinkedInConversionsClient: 

28 def __init__( 

29 self, 

30 *, 

31 api_version: str, 

32 client_id: str | None = None, 

33 client_secret: str | None = None, 

34 refresh_token: str | None = None, 

35 access_token: str | None = None, 

36 http: httpx.AsyncClient | None = None, 

37 ) -> None: 

38 self._client_id = client_id 

39 self._client_secret = client_secret 

40 self._api_version = api_version 

41 self._refresh_token = refresh_token 

42 self._access_token = access_token 

43 self._http = http or httpx.AsyncClient(timeout=15.0) 

44 

45 async def send(self, event: dict) -> None: 

46 response = await self._post(event, await self._token()) 

47 if response.status_code == 401 and self._refresh_token: 

48 response = await self._post(event, await self._token(force_refresh=True)) 

49 if response.status_code in (401, 403): 

50 # Fixed by an operator re-authorising the app, not by changing the event: hold the work. 

51 raise TransientDeliveryError( 

52 f"LinkedIn refused the access token ({response.status_code}); re-authorise the app" 

53 ) 

54 raise_for_response(response) 

55 

56 async def _post(self, event: dict, token: str) -> httpx.Response: 

57 try: 

58 return await self._http.post( 

59 CONVERSION_EVENTS_URL, 

60 json=event, 

61 headers={ 

62 "Authorization": f"Bearer {token}", 

63 "Linkedin-Version": self._api_version, 

64 "X-Restli-Protocol-Version": "2.0.0", 

65 }, 

66 ) 

67 except httpx.TransportError as exc: 

68 raise TransientDeliveryError(f"LinkedIn unreachable: {exc}") from exc 

69 

70 async def _token(self, *, force_refresh: bool = False) -> str: 

71 if self._access_token and not force_refresh: 

72 return self._access_token 

73 if not self._refresh_token or not self._client_id or not self._client_secret: 

74 raise TransientDeliveryError( 

75 "LinkedIn access token missing and no refresh token with client credentials configured" 

76 ) 

77 try: 

78 response = await self._http.post( 

79 TOKEN_URL, 

80 data={ 

81 "grant_type": "refresh_token", 

82 "refresh_token": self._refresh_token, 

83 "client_id": self._client_id, 

84 "client_secret": self._client_secret, 

85 }, 

86 ) 

87 except httpx.TransportError as exc: 

88 raise TransientDeliveryError(f"LinkedIn token endpoint unreachable: {exc}") from exc 

89 if not response.is_success: 

90 # invalid_grant arrives as 400: the refresh token was revoked, and only re-authorising fixes that. 

91 raise TransientDeliveryError(f"LinkedIn token refresh failed: {response.status_code} {response.text[:300]}") 

92 self._access_token = response.json()["access_token"] 

93 logger.info("LinkedIn access token refreshed") 

94 return self._access_token 

95 

96 async def close(self) -> None: 

97 await self._http.aclose()