Coverage for atlas_ecommerce/workers/provision.py: 72%

57 statements  

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

1"""Account provisioning for received orders. 

2 

3Runs out of band from the webhook so a slow or broken auth-proxy can never turn 

4into a non-2xx to Shopify, which after eight consecutive failures deletes the 

5subscription. Orders are retried until ``provision_max_attempts``, then parked as 

6failed for a human to look at. 

7""" 

8 

9import asyncio 

10from datetime import UTC, datetime, timedelta 

11from logging import getLogger 

12 

13from atlas_ecommerce.clients.auth_proxy import EcommerceAuthProxyClient, UserProvision 

14from atlas_ecommerce.config import Settings 

15from atlas_ecommerce.models import OrderDb, OrderStatus 

16 

17logger = getLogger(__name__) 

18 

19PROVENANCE = "shopify_preorder" 

20 

21# Exponential, capped. Without this a brief auth-proxy outage would burn every 

22# attempt within seconds and park orders that only needed a retry. 

23BACKOFF_BASE_SECONDS = 5 

24BACKOFF_CAP_SECONDS = 600 

25 

26 

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

28 """How long to wait before retrying an order that has failed this often.""" 

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

30 

31 

32async def provision_order(order: OrderDb, client: EcommerceAuthProxyClient, settings: Settings) -> None: 

33 """Create or link the Atlas account for one order.""" 

34 order.attempts += 1 

35 

36 if not order.email: 

37 # Nothing to provision against, and no retry will conjure an address. 

38 order.status = OrderStatus.FAILED 

39 order.last_error = "order has no email address" 

40 await order.save() 

41 logger.warning("Order %s carries no email; parked", order.shopify_order_id) 

42 return 

43 

44 try: 

45 user = await client.provision_user( 

46 UserProvision( 

47 email=order.email, 

48 first_name=order.first_name, 

49 last_name=order.last_name, 

50 trusted_metadata={ 

51 "provenance": PROVENANCE, 

52 "shopify_order_id": str(order.shopify_order_id), 

53 }, 

54 ) 

55 ) 

56 except Exception as exc: 

57 order.last_error = f"{type(exc).__name__}: {exc}" 

58 if order.attempts >= settings.provision_max_attempts: 

59 order.status = OrderStatus.FAILED 

60 logger.error("Order %s failed provisioning permanently", order.shopify_order_id) 

61 else: 

62 retry_in = backoff_for(order.attempts) 

63 order.next_attempt_at = datetime.now(UTC) + retry_in 

64 logger.warning( 

65 "Order %s provisioning attempt %s failed; retrying in %ss", 

66 order.shopify_order_id, 

67 order.attempts, 

68 int(retry_in.total_seconds()), 

69 ) 

70 await order.save() 

71 return 

72 

73 order.stytch_user_id = user.stytch_user_id 

74 order.status = OrderStatus.PROVISIONED 

75 order.provisioned_at = datetime.now(UTC) 

76 order.next_attempt_at = None 

77 order.last_error = None 

78 await order.save() 

79 logger.info("Order %s provisioned as %s", order.shopify_order_id, user.stytch_user_id) 

80 

81 

82async def drain_once(client: EcommerceAuthProxyClient, settings: Settings) -> int: 

83 """Provision one batch of pending orders. Returns how many were attempted.""" 

84 pending: list[OrderDb] = ( 

85 await OrderDb.find( 

86 { 

87 "status": OrderStatus.RECEIVED, 

88 "attempts": {"$lt": settings.provision_max_attempts}, 

89 "$or": [ 

90 {"next_attempt_at": None}, 

91 {"next_attempt_at": {"$lte": datetime.now(UTC)}}, 

92 ], 

93 }, 

94 projection_model=OrderDb, 

95 ) 

96 .limit(settings.provision_batch_size) 

97 .to_list() 

98 ) 

99 for order in pending: 

100 await provision_order(order, client, settings) 

101 return len(pending) 

102 

103 

104async def run_worker(client: EcommerceAuthProxyClient, settings: Settings) -> None: 

105 """Poll for received orders until cancelled.""" 

106 logger.info("Provisioning worker started") 

107 while True: 

108 try: 

109 attempted = await drain_once(client, settings) 

110 except asyncio.CancelledError: 

111 logger.info("Provisioning worker stopped") 

112 raise 

113 except Exception: 

114 logger.exception("Provisioning drain failed; continuing") 

115 attempted = 0 

116 if attempted == 0: 

117 await asyncio.sleep(settings.provision_poll_seconds)