Skip to content

Commit 557dfbe

Browse files
committed
feat(entitlements): lifecycle job for grace, lapse, reminders and reconcile
One scheduled tick a minute owns every transition no webhook fires: term end to grace, grace end to lapsed, the seven-day past-due ceiling, the 30/7/1-day reminders and the grace and lapsed emails (deduped on the audit log), and expired override pruning. Daily tasks compare live subscriptions with the provider and list pro features a switched-off flag withholds.
1 parent e154b4a commit 557dfbe

18 files changed

Lines changed: 1469 additions & 14 deletions

File tree

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
{
2+
"name": "Entitlements lifecycle job stalled",
3+
"description": "The lifecycle job ticks every minute and logs entitlements_lifecycle_tick. No tick for five minutes means subscriptions freeze in the generous direction: nobody moves to grace or lapses, reminders stop. Check the click worker (or the embedded scheduler) and the scheduled_tasks lease.",
4+
"type": "Threshold",
5+
"intervalMinutes": 5,
6+
"rangeMinutes": 5,
7+
"aplQuery": "['spoo-prod'] | where event == \"entitlements_lifecycle_tick\" | summarize count() by bin(_time, 5m)",
8+
"operator": "Below",
9+
"threshold": 1,
10+
"alertOnNoData": true,
11+
"notifierIds": [
12+
"xd9Sj0KGJ0e2TYPnl2"
13+
]
14+
}

dependencies/wiring.py

Lines changed: 50 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,7 @@
9393
from services.domain_intel_service import DomainIntelService
9494
from services.edge_cache.og_writethrough import OgEdgeWritethrough
9595
from services.entitlements import EntitlementService
96+
from services.entitlements.lifecycle import LifecycleService, lifecycle_tasks
9697
from services.entitlements.over_limit import OverLimitService
9798
from services.events.sinks import (
9899
InlineDomainEventSink,
@@ -330,6 +331,18 @@ def build_entitlement_service(
330331
)
331332

332333

334+
def build_billing_provider(settings: AppSettings, http_client) -> BillingProvider:
335+
billing = settings.billing
336+
if billing.provider == "paddle":
337+
return PaddleProvider(
338+
http_client,
339+
api_key=billing.paddle_api_key,
340+
webhook_secret=billing.paddle_webhook_secret,
341+
env=billing.paddle_env,
342+
)
343+
return NullBillingProvider()
344+
345+
333346
def build_billing_service(
334347
db,
335348
settings: AppSettings,
@@ -340,18 +353,8 @@ def build_billing_service(
340353
subscriptions: SubscriptionRepository,
341354
) -> BillingService:
342355
billing = settings.billing
343-
provider: BillingProvider
344-
if billing.provider == "paddle":
345-
provider = PaddleProvider(
346-
http_client,
347-
api_key=billing.paddle_api_key,
348-
webhook_secret=billing.paddle_webhook_secret,
349-
env=billing.paddle_env,
350-
)
351-
else:
352-
provider = NullBillingProvider()
353356
return BillingService(
354-
provider,
357+
build_billing_provider(settings, http_client),
355358
entitlements,
356359
subscriptions,
357360
BillingEventRepository(db["billing_events"]),
@@ -361,6 +364,32 @@ def build_billing_service(
361364
)
362365

363366

367+
def build_lifecycle_service(
368+
db,
369+
settings: AppSettings,
370+
*,
371+
store,
372+
entitlements: EntitlementService,
373+
provider: BillingProvider,
374+
mailer,
375+
) -> LifecycleService:
376+
"""``store`` is the caller's entitlement store tuple. An unconfigured
377+
mailer becomes None so the job never retries a send it cannot make."""
378+
_, events, subscriptions, overrides = store
379+
return LifecycleService(
380+
entitlements,
381+
subscriptions,
382+
overrides,
383+
events,
384+
UserRepository(db["users"]),
385+
CustomDomainRepository(db["custom_domains"]),
386+
FeatureFlagRepository(db["feature_flags"]),
387+
mailer if settings.email.zepto_api_token else None,
388+
app_url=settings.app_url,
389+
provider=provider,
390+
)
391+
392+
364393
def build_account_erasure_service(
365394
db,
366395
settings: AppSettings,
@@ -1178,6 +1207,16 @@ def wire_services(app: FastAPI, settings: AppSettings, redis_client) -> None:
11781207
),
11791208
),
11801209
erasure_sweep_task(account_erasure_service),
1210+
*lifecycle_tasks(
1211+
build_lifecycle_service(
1212+
db,
1213+
settings,
1214+
store=ent_store,
1215+
entitlements=app.state.entitlement_service,
1216+
provider=build_billing_provider(settings, http_client),
1217+
mailer=app.state.email_provider,
1218+
)
1219+
),
11811220
]
11821221
)
11831222
app.state.task_scheduler = TaskScheduler(

infrastructure/email/zeptomail.py

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -141,6 +141,20 @@ async def send_welcome_email(self, email: str, user_name: str | None) -> bool:
141141
)
142142
return await self._send(email, user_name, subject, html_body, text_body)
143143

144+
async def send_html(
145+
self,
146+
to_email: str,
147+
subject: str,
148+
template_name: str,
149+
context: dict,
150+
text_body: str,
151+
) -> bool:
152+
"""Render ``templates/emails/<template_name>`` with ``context`` and send it."""
153+
html_body = self._jinja.get_template(template_name).render(
154+
app_url=self._app_url, **context
155+
)
156+
return await self._send(to_email, None, subject, html_body, text_body)
157+
144158
async def send_deletion_requested(
145159
self, email: str, purge_after: datetime, restore_token: str | None = None
146160
) -> bool:

repositories/entitlement_event_repository.py

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@
22

33
from __future__ import annotations
44

5+
from datetime import datetime
6+
57
from bson import ObjectId
68

79
from repositories.base import BaseRepository
@@ -21,5 +23,27 @@ async def count_for(self, user_id: ObjectId) -> int:
2123
{"user_id": user_id, "kind": {"$nin": list(_VERSIONLESS_KINDS)}}
2224
)
2325

26+
async def has_reminder(
27+
self, user_id: ObjectId, kind: EntitlementEventKind, period: str
28+
) -> bool:
29+
doc = await self._find_one_raw(
30+
{"user_id": user_id, "kind": kind.value, "period": period}, {"_id": 1}
31+
)
32+
return doc is not None
33+
34+
async def users_with_override_writes_since(self, since: datetime) -> list[ObjectId]:
35+
return await self._col.distinct(
36+
"user_id",
37+
{
38+
"kind": {
39+
"$in": [
40+
EntitlementEventKind.OVERRIDE_GRANTED.value,
41+
EntitlementEventKind.OVERRIDE_REVOKED.value,
42+
]
43+
},
44+
"at": {"$gte": since},
45+
},
46+
)
47+
2448
async def delete_by_user(self, user_id: ObjectId) -> int:
2549
return await self._delete_many({"user_id": user_id})

repositories/entitlement_override_repository.py

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -153,6 +153,45 @@ async def revoke(
153153
await self._cache.invalidate(user_id)
154154
return True
155155

156+
async def find_expired(
157+
self, now: datetime, *, limit: int = 500
158+
) -> list[EntitlementOverrideDoc]:
159+
cursor = self._col.find({"expires_at": {"$ne": None, "$lte": now}}).limit(limit)
160+
docs = await cursor.to_list(length=limit)
161+
return [EntitlementOverrideDoc.from_mongo(d) for d in docs] # type: ignore[misc]
162+
163+
async def revoke_expired(
164+
self,
165+
user_id: ObjectId,
166+
key: str,
167+
now: datetime,
168+
*,
169+
actor: str,
170+
reason: str,
171+
) -> bool:
172+
"""Revoke only while still expired: an extension granted in between wins."""
173+
existing = await self._find_one({"user_id": user_id, "key": key})
174+
if existing is None or existing.expires_at is None:
175+
return False
176+
deleted = await self._delete(
177+
{"_id": existing.id, "expires_at": {"$ne": None, "$lte": now}}
178+
)
179+
if not deleted:
180+
return False
181+
await self._events.append(
182+
EntitlementEventDoc(
183+
user_id=user_id,
184+
kind=EntitlementEventKind.OVERRIDE_REVOKED,
185+
actor=actor,
186+
reason=reason,
187+
before=_snapshot(existing),
188+
after=None,
189+
at=now,
190+
)
191+
)
192+
await self._cache.invalidate(user_id)
193+
return True
194+
156195
async def delete_by_user(self, user_id: ObjectId) -> int:
157196
deleted = await self._delete_many({"user_id": user_id})
158197
if deleted:

repositories/indexes.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -199,6 +199,8 @@ async def ensure_indexes(
199199
await overrides_col.create_index([("expires_at", 1)], sparse=True)
200200
ent_events_col = db["entitlement_events"]
201201
await ent_events_col.create_index([("user_id", 1), ("at", -1)])
202+
# The lifecycle tick sweeps recent override writes by kind and time.
203+
await ent_events_col.create_index([("kind", 1), ("at", -1)])
202204
# The dedupe: a webhook delivery is handled once because this insert
203205
# either lands or raises DuplicateKeyError.
204206
billing_events_col = db["billing_events"]

repositories/subscription_repository.py

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -145,6 +145,84 @@ async def write_transition(
145145
await self._cache.invalidate(user_id)
146146
return after
147147

148+
async def find_term_ended(self, now: datetime) -> list[SubscriptionDoc]:
149+
"""Paid terms that have run out: a scheduled cancel past its period
150+
end, or a prepaid year past its end."""
151+
return await self._find_many(
152+
{
153+
"$or": [
154+
{
155+
"status": SubscriptionStatus.CANCEL_AT_PERIOD_END.value,
156+
"current_period_end": {"$lte": now},
157+
},
158+
{
159+
"status": SubscriptionStatus.ACTIVE.value,
160+
"kind": "prepaid",
161+
"prepaid_until": {"$lte": now},
162+
},
163+
]
164+
}
165+
)
166+
167+
async def find_grace_ended(self, now: datetime) -> list[SubscriptionDoc]:
168+
return await self._find_many(
169+
{"status": SubscriptionStatus.GRACE.value, "grace_until": {"$lte": now}}
170+
)
171+
172+
async def find_in_grace(self) -> list[SubscriptionDoc]:
173+
return await self._find_many({"status": SubscriptionStatus.GRACE.value})
174+
175+
async def find_lapsed_since(
176+
self, since: datetime, *, not_by: str
177+
) -> list[SubscriptionDoc]:
178+
"""Lapsed recently enough to still be owed the lapsed email, except
179+
those lapsed by ``not_by``; a lapsed document's last write is the lapse."""
180+
return await self._find_many(
181+
{
182+
"status": SubscriptionStatus.LAPSED.value,
183+
"updated_at": {"$gte": since},
184+
"lapsed_by": {"$ne": not_by},
185+
}
186+
)
187+
188+
async def find_past_due_older_than(self, cutoff: datetime) -> list[SubscriptionDoc]:
189+
return await self._find_many(
190+
{
191+
"status": SubscriptionStatus.PAST_DUE.value,
192+
"current_period_end": {"$lte": cutoff},
193+
}
194+
)
195+
196+
async def find_terms_ending_before(self, until: datetime) -> list[SubscriptionDoc]:
197+
"""Subscriptions whose paid term ends before ``until`` and will not
198+
renew: prepaid years and scheduled cancels."""
199+
return await self._find_many(
200+
{
201+
"$or": [
202+
{
203+
"status": SubscriptionStatus.CANCEL_AT_PERIOD_END.value,
204+
"current_period_end": {"$lte": until},
205+
},
206+
{
207+
"status": SubscriptionStatus.ACTIVE.value,
208+
"kind": "prepaid",
209+
"prepaid_until": {"$lte": until},
210+
},
211+
]
212+
}
213+
)
214+
215+
async def find_provider_live(self, provider: str) -> list[SubscriptionDoc]:
216+
"""Every non-lapsed recurring subscription of a provider, for reconcile."""
217+
return await self._find_many(
218+
{
219+
"provider": provider,
220+
"kind": "recurring",
221+
"status": {"$ne": SubscriptionStatus.LAPSED.value},
222+
},
223+
limit=5000,
224+
)
225+
148226
async def find_by_provider_subscription(
149227
self, subscription_id: str
150228
) -> SubscriptionDoc | None:

schemas/models/entitlement_event.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,4 +32,6 @@ class EntitlementEventDoc(MongoBaseModel):
3232
reason: str
3333
before: dict[str, Any] | None = None
3434
after: dict[str, Any] | None = None
35+
# Dedupe key for reminders: one email per (user, kind, period).
36+
period: str | None = None
3537
at: datetime

schemas/models/subscription.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,9 @@ class SubscriptionDoc(MongoBaseModel):
5757
grace_until: datetime | None = None
5858
founding: bool = False
5959
founding_streak_ok: bool = False
60+
# The state-machine event that lapsed it; "ended" means a refund or a
61+
# manual end, not a term running out.
62+
lapsed_by: str | None = None
6063
created_at: datetime | None = None
6164
updated_at: datetime | None = None
6265

0 commit comments

Comments
 (0)