Coverage for polar/customer/tasks.py: 37%
40 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 12:42 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 12:42 +0000
1import uuid
2from typing import Literal
4from sqlalchemy.orm import joinedload
6from polar.event.service import event as event_service
7from polar.event.system import CustomerUpdatedFields, SystemEvent, build_system_event
8from polar.exceptions import PolarTaskError
9from polar.models import Customer
10from polar.models.webhook_endpoint import CustomerWebhookEventType
11from polar.worker import AsyncSessionMaker, RedisMiddleware, TaskPriority, actor
13from .repository import CustomerRepository
14from .service import customer as customer_service
17class CustomerTaskError(PolarTaskError): ... 17 ↛ 20line 17 didn't jump to line 20 because
20class CustomerDoesNotExist(CustomerTaskError):
21 def __init__(self, customer_id: uuid.UUID) -> None:
22 self.customer_id = customer_id
23 message = f"The customer with id {customer_id} does not exist."
24 super().__init__(message)
27@actor(actor_name="customer.webhook", priority=TaskPriority.MEDIUM)
28async def customer_webhook(
29 event_type: CustomerWebhookEventType, customer_id: uuid.UUID
30) -> None:
31 async with AsyncSessionMaker() as session:
32 repository = CustomerRepository.from_session(session)
33 customer = await repository.get_by_id(
34 customer_id,
35 include_deleted=True,
36 options=(joinedload(Customer.organization),),
37 )
39 if customer is None:
40 raise CustomerDoesNotExist(customer_id)
42 await customer_service.webhook(
43 session, RedisMiddleware.get(), event_type, customer
44 )
47@actor(actor_name="customer.event", priority=TaskPriority.LOW)
48async def customer_event(
49 customer_id: uuid.UUID,
50 event_name: Literal[
51 SystemEvent.customer_created,
52 SystemEvent.customer_updated,
53 SystemEvent.customer_deleted,
54 ],
55 updated_fields: CustomerUpdatedFields | None = None,
56) -> None:
57 async with AsyncSessionMaker() as session:
58 repository = CustomerRepository.from_session(session)
59 customer = await repository.get_by_id(
60 customer_id,
61 include_deleted=True,
62 options=(joinedload(Customer.organization),),
63 )
65 if customer is None:
66 raise CustomerDoesNotExist(customer_id)
68 match event_name:
69 case SystemEvent.customer_created:
70 event = build_system_event(
71 event_name,
72 customer=customer,
73 organization=customer.organization,
74 metadata={
75 "customer_id": str(customer.id),
76 "customer_email": customer.email,
77 "customer_name": customer.name,
78 "customer_external_id": customer.external_id,
79 },
80 )
81 case SystemEvent.customer_deleted:
82 event = build_system_event(
83 event_name,
84 customer=customer,
85 organization=customer.organization,
86 metadata={
87 "customer_id": str(customer.id),
88 "customer_email": customer.email,
89 "customer_name": customer.name,
90 "customer_external_id": customer.external_id,
91 },
92 )
93 case SystemEvent.customer_updated:
94 event = build_system_event(
95 event_name,
96 customer=customer,
97 organization=customer.organization,
98 metadata={
99 "customer_id": str(customer.id),
100 "customer_email": customer.email,
101 "customer_name": customer.name,
102 "customer_external_id": customer.external_id,
103 "updated_fields": updated_fields or {},
104 },
105 )
107 await event_service.create_event(session, event)