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
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-25 16:36 +0000
1from datetime import timedelta
2from logging import getLogger
4import httpx
6from atlas_ecommerce.clients.delivery import TransientDeliveryError, raise_for_response
7from atlas_ecommerce.purchases import PurchaseEvent
9logger = getLogger(__name__)
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)
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 }
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)
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)
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
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
96 async def close(self) -> None:
97 await self._http.aclose()