Skip to content

Commit 1cc682f

Browse files
committed
fix(membership): order Apple subscription events
1 parent d3e4d03 commit 1cc682f

10 files changed

Lines changed: 322 additions & 4 deletions

apps/core/src/modules/membership/billing-webhook-event.repository.ts

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import { Inject, Injectable } from '@nestjs/common'
2-
import { and, eq } from 'drizzle-orm'
2+
import { and, asc, desc, eq, isNotNull, isNull, sql } from 'drizzle-orm'
33

44
import { PG_DB_TOKEN } from '~/constants/system.constant'
55
import { billingWebhookEvents } from '~/database/schema'
@@ -76,6 +76,53 @@ export class BillingWebhookEventRepository extends BaseRepository {
7676
return row ? mapRow(row) : null
7777
}
7878

79+
async findPendingByProviderSubscriptionId(
80+
provider: string,
81+
providerSubscriptionId: string,
82+
): Promise<BillingWebhookEventRow[]> {
83+
const rows = await this.db
84+
.select()
85+
.from(billingWebhookEvents)
86+
.where(
87+
and(
88+
eq(billingWebhookEvents.provider, provider),
89+
isNull(billingWebhookEvents.processedAt),
90+
sql`${billingWebhookEvents.payload}->'_normalizedMembershipEvent'->>'subscriptionId' = ${providerSubscriptionId}`,
91+
)!,
92+
)
93+
.orderBy(
94+
asc(
95+
sql`${billingWebhookEvents.payload}->'_normalizedMembershipEvent'->>'occurredAt'`,
96+
),
97+
asc(billingWebhookEvents.receivedAt),
98+
)
99+
return rows.map(mapRow)
100+
}
101+
102+
async findLatestProcessedByProviderSubscriptionId(
103+
provider: string,
104+
providerSubscriptionId: string,
105+
): Promise<BillingWebhookEventRow | null> {
106+
const [row] = await this.db
107+
.select()
108+
.from(billingWebhookEvents)
109+
.where(
110+
and(
111+
eq(billingWebhookEvents.provider, provider),
112+
isNotNull(billingWebhookEvents.processedAt),
113+
sql`${billingWebhookEvents.payload}->'_normalizedMembershipEvent'->>'subscriptionId' = ${providerSubscriptionId}`,
114+
)!,
115+
)
116+
.orderBy(
117+
desc(
118+
sql`${billingWebhookEvents.payload}->'_normalizedMembershipEvent'->>'occurredAt'`,
119+
),
120+
desc(billingWebhookEvents.receivedAt),
121+
)
122+
.limit(1)
123+
return row ? mapRow(row) : null
124+
}
125+
79126
async markProcessed(
80127
id: EntityId | string,
81128
processedAt: Date,

apps/core/src/modules/membership/membership.controller.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -238,6 +238,9 @@ export class MembershipController {
238238
verified.event.subscriptionId,
239239
)
240240
if (!bound) {
241+
if (verified.event.provider === 'apple') {
242+
await this.membershipService.deferEvent(verified)
243+
}
241244
return {
242245
ok: true,
243246
applied: false,

apps/core/src/modules/membership/membership.service.ts

Lines changed: 118 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,70 @@ import type {
2222
VerifiedBillingEvent,
2323
} from './providers/provider.interface'
2424

25+
const NORMALIZED_EVENT_PAYLOAD_KEY = '_normalizedMembershipEvent'
26+
27+
type StoredNormalizedBillingEvent = Omit<
28+
NormalizedBillingEvent,
29+
'currentPeriodEnd' | 'occurredAt'
30+
> & {
31+
currentPeriodEnd: string
32+
occurredAt?: string
33+
}
34+
35+
const storeBillingEventPayload = (
36+
event: NormalizedBillingEvent,
37+
rawPayload: unknown,
38+
): Record<string, unknown> => {
39+
const payload =
40+
rawPayload && typeof rawPayload === 'object' && !Array.isArray(rawPayload)
41+
? { ...(rawPayload as Record<string, unknown>) }
42+
: { rawPayload }
43+
const storedEvent: StoredNormalizedBillingEvent = {
44+
...event,
45+
currentPeriodEnd: event.currentPeriodEnd.toISOString(),
46+
occurredAt: event.occurredAt?.toISOString(),
47+
}
48+
payload[NORMALIZED_EVENT_PAYLOAD_KEY] = storedEvent
49+
return payload
50+
}
51+
52+
const readStoredBillingEvent = (
53+
payload: unknown,
54+
): NormalizedBillingEvent | null => {
55+
if (!payload || typeof payload !== 'object') return null
56+
const stored = (payload as Record<string, unknown>)[
57+
NORMALIZED_EVENT_PAYLOAD_KEY
58+
] as Partial<StoredNormalizedBillingEvent> | undefined
59+
if (
60+
!stored ||
61+
typeof stored.eventId !== 'string' ||
62+
typeof stored.provider !== 'string' ||
63+
typeof stored.type !== 'string' ||
64+
typeof stored.customerId !== 'string' ||
65+
typeof stored.subscriptionId !== 'string' ||
66+
typeof stored.currentPeriodEnd !== 'string' ||
67+
typeof stored.readerId !== 'string'
68+
) {
69+
return null
70+
}
71+
const currentPeriodEnd = new Date(stored.currentPeriodEnd)
72+
const occurredAt = stored.occurredAt ? new Date(stored.occurredAt) : undefined
73+
if (
74+
Number.isNaN(currentPeriodEnd.getTime()) ||
75+
(occurredAt && Number.isNaN(occurredAt.getTime()))
76+
) {
77+
return null
78+
}
79+
return {
80+
...(stored as Omit<
81+
NormalizedBillingEvent,
82+
'currentPeriodEnd' | 'occurredAt'
83+
>),
84+
currentPeriodEnd,
85+
occurredAt,
86+
}
87+
}
88+
2589
const isLiveProviderSubscription = (row: MembershipRow): boolean => {
2690
if (row.provider === 'manual') return false
2791
if (row.status !== 'active' && row.status !== 'on_hold') return false
@@ -115,6 +179,11 @@ export class MembershipService {
115179
rawPayload: input.decoded,
116180
rawType: 'apple.confirm',
117181
})
182+
await this.replayDeferredEvents(
183+
event.provider,
184+
event.subscriptionId,
185+
input.readerId,
186+
)
118187

119188
return this.toStatusResult(
120189
await this.membershipRepository.findByReaderId(input.readerId),
@@ -133,7 +202,7 @@ export class MembershipService {
133202
provider: event.provider,
134203
eventId: event.eventId,
135204
type: rawType,
136-
payload: rawPayload,
205+
payload: storeBillingEventPayload(event, rawPayload),
137206
})
138207

139208
if (!webhookEventRow) {
@@ -161,9 +230,57 @@ export class MembershipService {
161230
return { applied }
162231
}
163232

233+
async deferEvent(verifiedEvent: VerifiedBillingEvent): Promise<void> {
234+
const { event, rawPayload, rawType } = verifiedEvent
235+
await this.billingWebhookEventRepository.create({
236+
provider: event.provider,
237+
eventId: event.eventId,
238+
type: rawType,
239+
payload: storeBillingEventPayload(event, rawPayload),
240+
})
241+
}
242+
243+
private async replayDeferredEvents(
244+
provider: string,
245+
subscriptionId: string,
246+
readerId: string,
247+
): Promise<void> {
248+
const rows =
249+
await this.billingWebhookEventRepository.findPendingByProviderSubscriptionId(
250+
provider,
251+
subscriptionId,
252+
)
253+
for (const row of rows) {
254+
const storedEvent = readStoredBillingEvent(row.payload)
255+
if (storedEvent) {
256+
await this.applyMembershipState({ ...storedEvent, readerId })
257+
}
258+
await this.billingWebhookEventRepository.markProcessed(row.id, new Date())
259+
}
260+
}
261+
262+
private async isSupersededEvent(
263+
event: NormalizedBillingEvent,
264+
): Promise<boolean> {
265+
if (!event.occurredAt) return false
266+
const latestRow =
267+
await this.billingWebhookEventRepository.findLatestProcessedByProviderSubscriptionId(
268+
event.provider,
269+
event.subscriptionId,
270+
)
271+
if (!latestRow) return false
272+
const latestEvent = readStoredBillingEvent(latestRow.payload)
273+
return Boolean(
274+
latestEvent?.occurredAt &&
275+
latestEvent.occurredAt.getTime() >= event.occurredAt.getTime(),
276+
)
277+
}
278+
164279
private async applyMembershipState(
165280
event: NormalizedBillingEvent,
166281
): Promise<boolean> {
282+
if (await this.isSupersededEvent(event)) return false
283+
167284
let existing = await this.membershipRepository.findByProviderSubscriptionId(
168285
event.subscriptionId,
169286
)

apps/core/src/modules/membership/providers/apple-transaction.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ export interface AppleDecodedTransaction {
1313
originalTransactionId: string
1414
productId: string
1515
revocationDate?: number
16+
signedDate: number
1617
transactionId: string
1718
}
1819

@@ -59,6 +60,7 @@ export function appleActivatedEvent(
5960
plan,
6061
currentPeriodEnd: new Date(decoded.expiresDate),
6162
readerId,
63+
occurredAt: new Date(decoded.signedDate),
6264
}
6365
}
6466

apps/core/src/modules/membership/providers/apple.provider.ts

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,8 @@ export class AppleProvider implements PaymentProviderAdapter {
5858
!decoded.transactionId ||
5959
!decoded.originalTransactionId ||
6060
!decoded.productId ||
61-
!decoded.expiresDate
61+
!decoded.expiresDate ||
62+
!decoded.signedDate
6263
) {
6364
throw createAppException(
6465
AppErrorCode.MEMBERSHIP_APPLE_TRANSACTION_INVALID,
@@ -70,6 +71,7 @@ export class AppleProvider implements PaymentProviderAdapter {
7071
originalTransactionId: decoded.originalTransactionId,
7172
productId: decoded.productId,
7273
revocationDate: decoded.revocationDate,
74+
signedDate: decoded.signedDate,
7375
transactionId: decoded.transactionId,
7476
}
7577
} catch (error) {
@@ -119,6 +121,7 @@ export class AppleProvider implements PaymentProviderAdapter {
119121
let signedRenewalInfo: string | undefined
120122
let signedTransactionInfo: string | undefined
121123
let notificationUUID: string
124+
let notificationSignedDate: number | undefined
122125
try {
123126
const notification = await this.verifyNotificationWithFallback(
124127
signedPayload,
@@ -128,6 +131,7 @@ export class AppleProvider implements PaymentProviderAdapter {
128131
notificationType = notification.notificationType ?? ''
129132
notificationSubtype = notification.subtype ?? ''
130133
notificationUUID = notification.notificationUUID ?? notificationType
134+
notificationSignedDate = notification.signedDate
131135
signedRenewalInfo = notification.data?.signedRenewalInfo
132136
signedTransactionInfo = notification.data?.signedTransactionInfo
133137
} catch (error) {
@@ -161,7 +165,12 @@ export class AppleProvider implements PaymentProviderAdapter {
161165
bundleId,
162166
membershipConfig.appleAppAppleId,
163167
).catch(() => null)
164-
if (!decoded?.originalTransactionId || !decoded.expiresDate) {
168+
const occurredAt = notificationSignedDate ?? decoded?.signedDate
169+
if (
170+
!decoded?.originalTransactionId ||
171+
!decoded.expiresDate ||
172+
!Number.isFinite(occurredAt)
173+
) {
165174
return {
166175
ignored: true,
167176
rawType: notificationType,
@@ -216,6 +225,7 @@ export class AppleProvider implements PaymentProviderAdapter {
216225
plan,
217226
currentPeriodEnd: new Date(currentPeriodEnd),
218227
readerId: '',
228+
occurredAt: new Date(occurredAt!),
219229
},
220230
rawType: notificationType,
221231
rawPayload: {

apps/core/src/modules/membership/providers/provider.interface.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ export interface NormalizedBillingEvent {
99
plan?: MembershipPlan
1010
currentPeriodEnd: Date
1111
readerId: string
12+
occurredAt?: Date
1213
}
1314

1415
export interface VerifiedBillingEvent {

apps/core/test/src/modules/membership/apple-transaction.spec.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@ describe('appleActivatedEvent', () => {
4646
expiresDate: Date.parse('2026-09-01T00:00:00.000Z'),
4747
originalTransactionId: 'orig-1',
4848
productId: 'yohaku.membership.yearly',
49+
signedDate: Date.parse('2026-08-01T00:00:00.000Z'),
4950
transactionId: 'txn-1',
5051
},
5152
'reader-1',
@@ -63,6 +64,7 @@ describe('appleActivatedEvent', () => {
6364
expect(event.currentPeriodEnd.toISOString()).toBe(
6465
'2026-09-01T00:00:00.000Z',
6566
)
67+
expect(event.occurredAt?.toISOString()).toBe('2026-08-01T00:00:00.000Z')
6668
})
6769
})
6870

apps/core/test/src/modules/membership/apple.provider.spec.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ describe('AppleProvider', () => {
3636
originalTransactionId: 'original-1',
3737
productId: 'yohaku.membership.monthly',
3838
revocationDate: Date.parse('2026-08-21T00:00:00.000Z'),
39+
signedDate: Date.parse('2026-08-01T00:00:00.000Z'),
3940
transactionId: 'transaction-1',
4041
})
4142

@@ -68,6 +69,7 @@ describe('AppleProvider', () => {
6869
data: { signedTransactionInfo: 'signed-transaction' },
6970
notificationType,
7071
notificationUUID: 'notification-1',
72+
signedDate: Date.parse('2026-08-21T00:00:00.000Z'),
7173
})
7274
vi.spyOn(
7375
provider as any,
@@ -107,6 +109,7 @@ describe('AppleProvider', () => {
107109
},
108110
notificationType: 'DID_FAIL_TO_RENEW',
109111
notificationUUID: 'notification-1',
112+
signedDate: Date.parse('2026-08-21T00:00:00.000Z'),
110113
subtype: 'GRACE_PERIOD',
111114
})
112115
vi.spyOn(
@@ -133,6 +136,7 @@ describe('AppleProvider', () => {
133136
expect(result).toMatchObject({
134137
event: {
135138
currentPeriodEnd: new Date('2026-08-27T00:00:00.000Z'),
139+
occurredAt: new Date('2026-08-21T00:00:00.000Z'),
136140
type: 'on_hold',
137141
},
138142
})

apps/core/test/src/modules/membership/membership.controller.e2e-spec.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -594,6 +594,7 @@ describe('MembershipController (e2e)', () => {
594594
expiresDate: Date.now() + 86_400_000,
595595
originalTransactionId: 'orig-e2e',
596596
productId: 'yohaku.membership.monthly',
597+
signedDate: Date.now(),
597598
transactionId: 'txn-e2e',
598599
})
599600

@@ -628,6 +629,7 @@ describe('MembershipController (e2e)', () => {
628629
expiresDate: Date.now() + 86_400_000,
629630
originalTransactionId: 'orig-other-reader',
630631
productId: 'yohaku.membership.monthly',
632+
signedDate: Date.now(),
631633
transactionId: 'txn-other-reader',
632634
})
633635

0 commit comments

Comments
 (0)