Coverage for atlas_ecommerce/workers/conversions.py: 94%
230 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
1import asyncio
2from collections.abc import Awaitable
3from dataclasses import dataclass
4from datetime import UTC, datetime, timedelta
5from logging import getLogger
6from uuid import uuid4
8from pymongo import ReturnDocument
9from pymongo.errors import DuplicateKeyError
11from atlas_ecommerce.clients import google_ads, linkedin_conversions, meta_capi, mixpanel, x_conversions
12from atlas_ecommerce.clients.delivery import PermanentDeliveryError, TransientDeliveryError
13from atlas_ecommerce.clients.google_ads import GoogleDataManagerClient
14from atlas_ecommerce.clients.linkedin_conversions import LinkedInConversionsClient
15from atlas_ecommerce.clients.meta_capi import MetaConversionsClient
16from atlas_ecommerce.clients.mixpanel import MixpanelClient
17from atlas_ecommerce.clients.x_conversions import XConversionsClient
18from atlas_ecommerce.config import Secrets, Settings
19from atlas_ecommerce.models import ConversionDispatchDb, Destination, DispatchStatus
20from atlas_ecommerce.purchases import PurchaseEvent, PurchaseSource
22logger = getLogger(__name__)
24# Each vendor refuses events older than its window; a purchase past it is skipped, never sent or failed.
25MAX_EVENT_AGE: dict[Destination, timedelta] = {
26 Destination.LINKEDIN: linkedin_conversions.MAX_EVENT_AGE,
27 Destination.X: x_conversions.MAX_EVENT_AGE,
28 Destination.MIXPANEL: mixpanel.MAX_EVENT_AGE,
29 Destination.META: meta_capi.MAX_EVENT_AGE,
30 Destination.GOOGLE_ADS: google_ads.MAX_EVENT_AGE,
31}
33# Doubling from 30 s and capped at an hour, so the default attempt budget spans about a day of vendor
34# outage or expired credentials instead of the minutes the provisioning schedule allows.
35BACKOFF_BASE_SECONDS = 30
36BACKOFF_CAP_SECONDS = 3600
39def backoff_for(attempts: int) -> timedelta:
40 return timedelta(seconds=min(BACKOFF_BASE_SECONDS * 2 ** (attempts - 1), BACKOFF_CAP_SECONDS))
43def utcnow() -> datetime:
44 return datetime.now(UTC)
47@dataclass
48class ConversionClients:
49 linkedin: LinkedInConversionsClient | None = None
50 x: XConversionsClient | None = None
51 mixpanel: MixpanelClient | None = None
52 meta: MetaConversionsClient | None = None
53 google_ads: GoogleDataManagerClient | None = None
55 def configured(self, destination: Destination) -> bool:
56 return getattr(self, destination.value) is not None
58 def destinations(self) -> list[Destination]:
59 return [destination for destination in Destination if self.configured(destination)]
61 async def send(self, destination: Destination, payload: dict) -> None:
62 client = getattr(self, destination.value)
63 if client is None:
64 raise TransientDeliveryError(f"{destination.value} is not configured on this replica")
65 await client.send(payload)
67 async def close(self) -> None:
68 for client in (self.linkedin, self.x, self.mixpanel, self.meta, self.google_ads):
69 if client is not None:
70 await client.close()
73def _secret(value) -> str | None:
74 return value.get_secret_value() if value is not None else None
77def build_clients(settings: Settings, secrets: Secrets) -> ConversionClients:
78 clients = ConversionClients()
80 linkedin_id = _secret(secrets.linkedin_client_id)
81 linkedin_secret = _secret(secrets.linkedin_client_secret)
82 linkedin_refresh = _secret(secrets.linkedin_refresh_token)
83 linkedin_access = _secret(secrets.linkedin_access_token)
84 refreshable = bool(linkedin_refresh and linkedin_id and linkedin_secret)
85 if settings.linkedin_conversion_urn and (linkedin_access or refreshable):
86 clients.linkedin = LinkedInConversionsClient(
87 api_version=settings.linkedin_api_version,
88 client_id=linkedin_id,
89 client_secret=linkedin_secret,
90 refresh_token=linkedin_refresh if refreshable else None,
91 access_token=linkedin_access,
92 )
94 x_token = _secret(secrets.x_pixel_token)
95 if settings.x_pixel_id and settings.x_purchase_event_id and x_token:
96 clients.x = XConversionsClient(pixel_id=settings.x_pixel_id, pixel_token=x_token)
98 if settings.mixpanel_token:
99 clients.mixpanel = MixpanelClient()
101 meta_token = _secret(secrets.meta_capi_token)
102 if settings.meta_pixel_id and meta_token:
103 clients.meta = MetaConversionsClient(
104 pixel_id=settings.meta_pixel_id, access_token=meta_token, api_version=settings.meta_api_version
105 )
107 ads_client_id = _secret(secrets.google_ads_client_id)
108 ads_client_secret = _secret(secrets.google_ads_client_secret)
109 ads_refresh_token = _secret(secrets.google_ads_refresh_token)
110 if (
111 settings.google_ads_customer_id
112 and settings.google_ads_purchase_conversion_action
113 and ads_client_id
114 and ads_client_secret
115 and ads_refresh_token
116 ):
117 clients.google_ads = GoogleDataManagerClient(
118 client_id=ads_client_id, client_secret=ads_client_secret, refresh_token=ads_refresh_token
119 )
121 for destination in Destination:
122 state = "enabled" if clients.configured(destination) else "waiting for configuration"
123 logger.info("Conversion destination %s %s", destination.value, state)
124 return clients
127def skip_reason(destination: Destination, purchase: PurchaseEvent, now: datetime | None = None) -> str | None:
128 attribution = purchase.attribution
129 if purchase.problem:
130 return purchase.problem
131 if not purchase.paid:
132 return "order not paid"
133 if not attribution.present:
134 return "order carries no attribution"
135 if destination is Destination.MIXPANEL:
136 if not attribution.analytics:
137 return "no analytics consent"
138 if not attribution.mp_distinct_id:
139 return "no mixpanel distinct_id"
140 else:
141 if not attribution.marketing:
142 return "no marketing consent"
143 if attribution.gpc:
144 return "global privacy control opt-out"
145 match destination:
146 case Destination.META:
147 if not attribution.user_agent:
148 return "no user agent"
149 if not attribution.fbp and not attribution.fbc and "fbclid" not in attribution.click_ids:
150 return "no fbp, fbc or fbclid"
151 case Destination.GOOGLE_ADS:
152 if not any(key in attribution.click_ids for key in ("gclid", "gbraid", "wbraid")):
153 return "no gclid, gbraid or wbraid"
154 case _:
155 click_id = "li_fat_id" if destination is Destination.LINKEDIN else "twclid"
156 if click_id not in attribution.click_ids:
157 return f"no {click_id}"
158 window = MAX_EVENT_AGE[destination]
159 if (now or utcnow()) - purchase.happened_at > window:
160 return f"purchase older than {destination.value}'s {window.days}-day window"
161 return None
164def plan(purchase: PurchaseEvent) -> list[ConversionDispatchDb]:
165 dispatches = []
166 for destination in Destination:
167 reason = skip_reason(destination, purchase)
168 dispatches.append(
169 ConversionDispatchDb(
170 provider=purchase.provider,
171 order_id=purchase.order_id,
172 destination=destination,
173 status=DispatchStatus.SKIPPED if reason else DispatchStatus.PENDING,
174 skip_reason=reason,
175 )
176 )
177 return dispatches
180def build_payload(destination: Destination, purchase: PurchaseEvent, settings: Settings) -> dict:
181 match destination:
182 case Destination.LINKEDIN:
183 return linkedin_conversions.purchase_event(purchase, settings.linkedin_conversion_urn)
184 case Destination.X:
185 return x_conversions.purchase_conversion(purchase, settings.x_purchase_event_id)
186 case Destination.MIXPANEL:
187 return mixpanel.purchase_event(purchase, settings.mixpanel_token)
188 case Destination.META:
189 return meta_capi.purchase_event(purchase, settings.meta_test_event_code or None)
190 case Destination.GOOGLE_ADS:
191 return google_ads.ingest_request(
192 purchase,
193 customer_id=settings.google_ads_customer_id,
194 conversion_action=settings.google_ads_purchase_conversion_action,
195 login_customer_id=settings.google_ads_login_customer_id or None,
196 validate_only=settings.google_data_manager_validate_only,
197 )
200async def enqueue_once(sources: list[PurchaseSource], settings: Settings) -> int:
201 """Create the dispatches for purchases that have none yet. Returns how many purchases were planned."""
202 planned = 0
203 for source in sources:
204 for purchase in await source.unplanned(settings.conversion_batch_size):
205 # The dispatches are the record of planning; the source's marker is only a fast-path
206 # filter that another writer's snapshot save may reset. Inserting just the missing
207 # destinations makes a re-plan cheap and completes a batch a crash left half-written.
208 existing = await planned_destinations(purchase)
209 inserted = 0
210 for dispatch in plan(purchase):
211 if dispatch.destination in existing:
212 continue
213 try:
214 await dispatch.insert()
215 inserted += 1
216 except DuplicateKeyError:
217 pass
218 await source.mark_planned(purchase)
219 if inserted:
220 planned += 1
221 logger.info("%s order %s planned for conversion delivery", purchase.provider, purchase.order_id)
222 return planned
225async def planned_destinations(purchase: PurchaseEvent) -> set[Destination]:
226 dispatches: list[ConversionDispatchDb] = await ConversionDispatchDb.find(
227 {"provider": purchase.provider, "order_id": purchase.order_id}, projection_model=ConversionDispatchDb
228 ).to_list()
229 return {dispatch.destination for dispatch in dispatches}
232async def claim_next(settings: Settings, destinations: list[Destination]) -> ConversionDispatchDb | None:
233 """Atomically lease one due dispatch this replica can deliver. The claim itself counts the attempt,
234 so a delivery that dies before reporting still consumes budget instead of being re-sent forever."""
235 if not destinations:
236 return None
237 now = utcnow()
238 due = {"status": DispatchStatus.PENDING, "attempts": {"$lt": settings.conversion_max_attempts}}
239 claimable = {
240 "destination": {"$in": [destination.value for destination in destinations]},
241 "$or": [
242 {**due, "next_attempt_at": None},
243 {**due, "next_attempt_at": {"$lte": now}},
244 # A lease older than the timeout belongs to a worker that died mid-delivery. Reclaimed
245 # whatever the count, so a dispatch that spent its last attempt that way still gets parked.
246 {
247 "status": DispatchStatus.PROCESSING,
248 "claimed_at": {"$lte": now - timedelta(seconds=settings.conversion_lease_seconds)},
249 },
250 ],
251 }
252 raw = await ConversionDispatchDb.get_pymongo_collection().find_one_and_update(
253 claimable,
254 {
255 "$set": {"status": DispatchStatus.PROCESSING, "claimed_at": now, "claim_token": uuid4().hex},
256 "$inc": {"attempts": 1},
257 },
258 return_document=ReturnDocument.AFTER,
259 )
260 return ConversionDispatchDb.model_validate(raw) if raw else None
263async def persist_outcome(dispatch: ConversionDispatchDb) -> bool:
264 """Write the result only if this worker still holds the lease; a stale worker's outcome is dropped."""
265 result = await ConversionDispatchDb.get_pymongo_collection().update_one(
266 {"_id": dispatch.id, "claim_token": dispatch.claim_token},
267 {
268 "$set": {
269 "status": dispatch.status,
270 "skip_reason": dispatch.skip_reason,
271 "last_error": dispatch.last_error,
272 "next_attempt_at": dispatch.next_attempt_at,
273 "sent_at": dispatch.sent_at,
274 "claimed_at": None,
275 "claim_token": None,
276 }
277 },
278 )
279 return result.modified_count == 1
282async def attempt(
283 dispatch: ConversionDispatchDb, sources: dict[str, PurchaseSource], clients: ConversionClients, settings: Settings
284) -> str | None:
285 """Send the conversion, or return the reason it must not be sent. Raises on delivery errors."""
286 if dispatch.attempts > settings.conversion_max_attempts:
287 raise PermanentDeliveryError("attempt budget spent by deliveries that never reported an outcome")
288 source = sources.get(dispatch.provider)
289 if source is None:
290 raise PermanentDeliveryError(f"no purchase source for provider {dispatch.provider!r}")
291 purchase = await source.load(dispatch.order_id)
292 if purchase is None:
293 raise PermanentDeliveryError("order record missing")
294 # Planning-time gates hold at delivery too: a dispatch that waited past a vendor's window is skipped.
295 reason = skip_reason(dispatch.destination, purchase)
296 if reason:
297 return reason
298 await clients.send(dispatch.destination, build_payload(dispatch.destination, purchase, settings))
299 return None
302async def deliver(
303 dispatch: ConversionDispatchDb, sources: dict[str, PurchaseSource], clients: ConversionClients, settings: Settings
304) -> None:
305 label = f"{dispatch.provider} order {dispatch.order_id}: {dispatch.destination}"
306 try:
307 reason = await attempt(dispatch, sources, clients, settings)
308 except PermanentDeliveryError as exc:
309 dispatch.status = DispatchStatus.FAILED
310 dispatch.last_error = f"{type(exc).__name__}: {exc}"
311 dispatch.next_attempt_at = None
312 logger.error("%s delivery failed permanently: %s", label, exc)
313 except Exception as exc:
314 dispatch.last_error = f"{type(exc).__name__}: {exc}"
315 if dispatch.attempts >= settings.conversion_max_attempts:
316 dispatch.status = DispatchStatus.FAILED
317 dispatch.next_attempt_at = None
318 logger.error("%s delivery gave up after %s attempts: %s", label, dispatch.attempts, exc)
319 else:
320 retry_in = backoff_for(dispatch.attempts)
321 dispatch.status = DispatchStatus.PENDING
322 dispatch.next_attempt_at = utcnow() + retry_in
323 logger.warning(
324 "%s delivery attempt %s failed; retrying in %ss: %s",
325 label,
326 dispatch.attempts,
327 int(retry_in.total_seconds()),
328 exc,
329 )
330 else:
331 dispatch.next_attempt_at = None
332 dispatch.last_error = None
333 if reason:
334 dispatch.status = DispatchStatus.SKIPPED
335 dispatch.skip_reason = reason
336 logger.info("%s conversion skipped at delivery: %s", label, reason)
337 else:
338 dispatch.status = DispatchStatus.SENT
339 dispatch.sent_at = utcnow()
340 logger.info("%s conversion delivered", label)
341 if not await persist_outcome(dispatch):
342 logger.warning("%s outcome discarded: the lease had moved to another worker", label)
345async def drain_once(sources: list[PurchaseSource], clients: ConversionClients, settings: Settings) -> int:
346 by_provider = {source.provider: source for source in sources}
347 delivered = 0
348 while delivered < settings.conversion_batch_size:
349 dispatch = await claim_next(settings, clients.destinations())
350 if dispatch is None:
351 break
352 await deliver(dispatch, by_provider, clients, settings)
353 delivered += 1
354 return delivered
357async def guarded(step: Awaitable[int], what: str) -> int:
358 try:
359 return await step
360 except Exception:
361 logger.exception("Conversion %s failed; continuing", what)
362 return 0
365async def run_worker(sources: list[PurchaseSource], clients: ConversionClients, settings: Settings) -> None:
366 logger.info("Conversion worker started")
367 while True:
368 try:
369 # Two guards, so a poison order in planning cannot starve the deliveries already queued.
370 worked = await guarded(enqueue_once(sources, settings), "planning")
371 worked += await guarded(drain_once(sources, clients, settings), "drain")
372 except asyncio.CancelledError:
373 logger.info("Conversion worker stopped")
374 raise
375 if worked == 0:
376 await asyncio.sleep(settings.conversion_poll_seconds)