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

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 

7 

8from pymongo import ReturnDocument 

9from pymongo.errors import DuplicateKeyError 

10 

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 

21 

22logger = getLogger(__name__) 

23 

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} 

32 

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 

37 

38 

39def backoff_for(attempts: int) -> timedelta: 

40 return timedelta(seconds=min(BACKOFF_BASE_SECONDS * 2 ** (attempts - 1), BACKOFF_CAP_SECONDS)) 

41 

42 

43def utcnow() -> datetime: 

44 return datetime.now(UTC) 

45 

46 

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 

54 

55 def configured(self, destination: Destination) -> bool: 

56 return getattr(self, destination.value) is not None 

57 

58 def destinations(self) -> list[Destination]: 

59 return [destination for destination in Destination if self.configured(destination)] 

60 

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) 

66 

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

71 

72 

73def _secret(value) -> str | None: 

74 return value.get_secret_value() if value is not None else None 

75 

76 

77def build_clients(settings: Settings, secrets: Secrets) -> ConversionClients: 

78 clients = ConversionClients() 

79 

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 ) 

93 

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) 

97 

98 if settings.mixpanel_token: 

99 clients.mixpanel = MixpanelClient() 

100 

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 ) 

106 

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 ) 

120 

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 

125 

126 

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 

162 

163 

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 

178 

179 

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 ) 

198 

199 

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 

223 

224 

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} 

230 

231 

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 

261 

262 

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 

280 

281 

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 

300 

301 

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) 

343 

344 

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 

355 

356 

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 

363 

364 

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)