Skip to content

Commit 07bcbfb

Browse files
committed
fix(billing): gate enterprise usage emails on an atomic claim and key the cache by source
The execution logger still summed an enterprise organization's whole reporting window on every execution completion, only to compute the baseline for the edge-triggered credits email. An edge needs an exact baseline, so the cached sum could not serve it. For reporting windows the organization credits email is now level-triggered. After the run's usage is recorded, the logger reads the window's usage through readSoftGateUsageCost and hands it to maybeSendReportingCreditsThresholdEmail. That function compares the level to a 'credits' claim on organization.limitNotifications, the same dedup the storage, tables and seats emails use: re-arm below 70%, claim 80 or 100 with one conditional UPDATE, and send only when this call won the claim. Opted-out admins are dropped before the claim, so they never burn it. The claim flow is now one shared notifyThresholdLevel in limit-notifications. It reads the notified threshold by primary key first, so a completion below the band with nothing armed, or above a threshold already sent, writes nothing and never resolves recipients. The email content, recipients, opt-outs and billing link are shared with the edge path through sendPaidCreditsEmail. A cached level can trail the ledger by up to the cache TTL, so a crossing is emailed up to one TTL late, on the first completion that sees it. It is never sent twice and never dropped while usage stays over the threshold. A new reporting window re-arms on its first completion because usage starts again from zero. Stripe and default organization periods and personal accounts keep the exact, edge-triggered path unchanged. The reporting-window cache is now reachable only through readSoftGateUsageCost. The raw cached reader is no longer exported and accepts only a reporting period, and the cache key includes the period's source, so another kind of period with the same bounds can never share a reporting sum.
1 parent 1757818 commit 07bcbfb

9 files changed

Lines changed: 811 additions & 180 deletions

File tree

‎apps/sim/lib/billing/core/limit-notifications.test.ts‎

Lines changed: 33 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,10 @@ vi.mock('@/components/emails', () => ({
3232
}))
3333
vi.mock('@sim/platform-authz/workspace', () => ({ isOrgAdminRole: isOrgAdminRoleMock }))
3434

35-
import { maybeSendLimitThresholdEmail } from '@/lib/billing/core/limit-notifications'
35+
import {
36+
maybeSendLimitThresholdEmail,
37+
notifyThresholdLevel,
38+
} from '@/lib/billing/core/limit-notifications'
3639

3740
const baseUserParams = {
3841
category: 'storage' as const,
@@ -141,12 +144,41 @@ describe('maybeSendLimitThresholdEmail', () => {
141144
})
142145

143146
it('re-arms but does not send when usage is fully cleared (zero usage)', async () => {
147+
queueTableRows(schemaMock.userStats, [{ notifications: { storage: 100 } }])
144148
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 0, limit: 5 })
145149
expect(dbChainMockFns.update).toHaveBeenCalledTimes(1)
146150
expect(dbChainMockFns.returning).not.toHaveBeenCalled()
147151
expect(sendEmailSpy).not.toHaveBeenCalled()
148152
})
149153

154+
it('writes nothing below the band when the category is not armed', async () => {
155+
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 0, limit: 5 })
156+
expect(dbChainMockFns.update).not.toHaveBeenCalled()
157+
})
158+
159+
it('skips recipients and the claim when the threshold was already notified', async () => {
160+
queueTableRows(schemaMock.userStats, [{ notifications: { storage: 80 } }])
161+
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4.5, limit: 5 })
162+
expect(getEmailPreferencesMock).not.toHaveBeenCalled()
163+
expect(dbChainMockFns.update).not.toHaveBeenCalled()
164+
expect(sendEmailSpy).not.toHaveBeenCalled()
165+
})
166+
167+
it('keys claims per category, so a notified storage threshold never blocks credits', async () => {
168+
queueTableRows(schemaMock.userStats, [{ notifications: { storage: 100 } }])
169+
const send = vi.fn(() => Promise.resolve())
170+
await notifyThresholdLevel({
171+
scope: 'user',
172+
stateId: 'u1',
173+
category: 'credits',
174+
percent: 100,
175+
resolveRecipients: () => Promise.resolve([{ email: 'u1@example.com' }]),
176+
send,
177+
})
178+
expect(dbChainMockFns.returning).toHaveBeenCalledTimes(1)
179+
expect(send).toHaveBeenCalledWith(100, [{ email: 'u1@example.com' }])
180+
})
181+
150182
it('skips when the limit is non-positive', async () => {
151183
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4, limit: 0 })
152184
expect(dbChainMockFns.update).not.toHaveBeenCalled()

‎apps/sim/lib/billing/core/limit-notifications.ts‎

Lines changed: 156 additions & 73 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { db } from '@sim/db'
22
import { member, organization, settings, user, userStats } from '@sim/db/schema'
33
import { createLogger } from '@sim/logger'
44
import { isOrgAdminRole } from '@sim/platform-authz/workspace'
5+
import { toNumberOrNull } from '@sim/utils/coerce'
56
import { and, eq, sql } from 'drizzle-orm'
67
import type { HighestPrioritySubscription } from '@/lib/billing/core/plan'
78
import { getHighestPrioritySubscription } from '@/lib/billing/core/subscription'
@@ -17,6 +18,18 @@ const logger = createLogger('LimitNotifications')
1718
/** Limit categories that send per-category threshold emails (credits has its own path). */
1819
export type LimitCategory = Extract<UpgradeReason, 'storage' | 'tables' | 'seats'>
1920

21+
/**
22+
* Every category whose threshold emails are deduplicated through `limitNotifications`: the limit
23+
* categories above, plus credits for the level-triggered enterprise reporting-window email in
24+
* `usage.ts`. Each category is its own key in the jsonb map, so their claims never interact.
25+
*/
26+
export type ThresholdCategory = LimitCategory | Extract<UpgradeReason, 'credits'>
27+
28+
/** A threshold a notification is sent at: the 80% warning or the 100% limit reached. */
29+
export type NotifiedThreshold = 80 | 100
30+
31+
type ThresholdScope = 'user' | 'organization'
32+
2033
const WARN_THRESHOLD = 80
2134
const REACH_THRESHOLD = 100
2235
/** Usage must drop below this band before the same threshold can re-notify (hysteresis). */
@@ -26,7 +39,7 @@ const REARM_BELOW = 70
2639
* Resolve the threshold a given usage percent should be notified at:
2740
* 100 at/over the limit, 80 when approaching, 0 otherwise.
2841
*/
29-
function thresholdFor(percent: number): 0 | 80 | 100 {
42+
function thresholdFor(percent: number): 0 | NotifiedThreshold {
3043
if (percent >= REACH_THRESHOLD) return REACH_THRESHOLD
3144
if (percent >= WARN_THRESHOLD) return WARN_THRESHOLD
3245
return 0
@@ -39,10 +52,10 @@ function thresholdFor(percent: number): 0 | 80 | 100 {
3952
* both claim, so the email is sent exactly once per crossing.
4053
*/
4154
async function claimThreshold(
42-
scope: 'user' | 'organization',
55+
scope: ThresholdScope,
4356
id: string,
44-
category: LimitCategory,
45-
threshold: number
57+
category: ThresholdCategory,
58+
threshold: NotifiedThreshold
4659
): Promise<boolean> {
4760
const setExpr = sql`jsonb_set(coalesce(${scope === 'user' ? userStats.limitNotifications : organization.limitNotifications}, '{}'::jsonb), ARRAY[${category}], to_jsonb(${threshold}::int))`
4861
const onlyIfLower =
@@ -68,9 +81,9 @@ async function claimThreshold(
6881

6982
/** Re-arm a category (reset its stored threshold to 0) once usage falls back into the low band. */
7083
async function rearmThreshold(
71-
scope: 'user' | 'organization',
84+
scope: ThresholdScope,
7285
id: string,
73-
category: LimitCategory
86+
category: ThresholdCategory
7487
): Promise<void> {
7588
const setExpr = sql`jsonb_set(coalesce(${scope === 'user' ? userStats.limitNotifications : organization.limitNotifications}, '{}'::jsonb), ARRAY[${category}], to_jsonb(0::int))`
7689
const onlyIfArmed =
@@ -91,7 +104,32 @@ async function rearmThreshold(
91104
}
92105
}
93106

94-
interface LimitEmailRecipient {
107+
/**
108+
* The threshold already notified for a category, read without taking a lock. It only lets a
109+
* caller skip work — the re-arm and claim UPDATEs stay conditional, so a read that races a
110+
* concurrent claim or re-arm can at worst cost one lost claim, which the next evaluation retries.
111+
*/
112+
async function readNotifiedThreshold(
113+
scope: ThresholdScope,
114+
id: string,
115+
category: ThresholdCategory
116+
): Promise<number> {
117+
const [row] =
118+
scope === 'user'
119+
? await db
120+
.select({ notifications: userStats.limitNotifications })
121+
.from(userStats)
122+
.where(eq(userStats.userId, id))
123+
.limit(1)
124+
: await db
125+
.select({ notifications: organization.limitNotifications })
126+
.from(organization)
127+
.where(eq(organization.id, id))
128+
.limit(1)
129+
return toNumberOrNull(row?.notifications[category]) ?? 0
130+
}
131+
132+
export interface LimitEmailRecipient {
95133
email: string
96134
name?: string
97135
}
@@ -108,8 +146,8 @@ async function isUnsubscribed(email: string): Promise<boolean> {
108146
* Returning an empty list means "nobody to notify" — the caller then skips the
109147
* claim so the dedup state isn't burned without an email going out.
110148
*/
111-
async function resolveRecipients(
112-
scope: 'user' | 'organization',
149+
export async function resolveLimitEmailRecipients(
150+
scope: ThresholdScope,
113151
params: { userId?: string; userEmail?: string; userName?: string; organizationId?: string }
114152
): Promise<LimitEmailRecipient[]> {
115153
if (scope === 'user') {
@@ -148,19 +186,62 @@ async function resolveRecipients(
148186
return recipients
149187
}
150188

189+
/**
190+
* Evaluate a threshold notification against the current usage level and send it at most once
191+
* per crossing. The shared dedup flow behind every threshold email in this module and the
192+
* credits email for enterprise reporting windows.
193+
*
194+
* The notified threshold is read first, so the steady state — usage below the band with nothing
195+
* armed, or above a threshold already notified — costs one primary-key read and never writes.
196+
* Below {@link REARM_BELOW} an armed category is re-armed; at or above a threshold not yet
197+
* notified, eligible recipients are resolved with opt-outs applied BEFORE the atomic
198+
* {@link claimThreshold}, so an opted-out recipient never burns the threshold (which would
199+
* suppress a later email once notifications are re-enabled). `send` runs only when this call won
200+
* the claim. Re-arm and claim are mutually exclusive per call, so the dedup stays a single atomic
201+
* claim with no re-arm/claim interleaving race.
202+
*
203+
* Because the decision depends only on the current level and the persisted claim, a caller may
204+
* pass a level that trails the truth: a crossing it misses is claimed on a later evaluation, and
205+
* a crossing it sees twice is claimed once.
206+
*/
207+
export async function notifyThresholdLevel(params: {
208+
scope: ThresholdScope
209+
/** The user id for `user` scope, the organization id for `organization` scope. */
210+
stateId: string
211+
category: ThresholdCategory
212+
/** Current usage as a percent of the limit. */
213+
percent: number
214+
/** Evaluate only the re-arm; never claim or send. */
215+
rearmOnly?: boolean
216+
resolveRecipients: () => Promise<LimitEmailRecipient[]>
217+
send: (threshold: NotifiedThreshold, recipients: LimitEmailRecipient[]) => Promise<void>
218+
}): Promise<void> {
219+
const { scope, stateId, category, percent } = params
220+
const notified = await readNotifiedThreshold(scope, stateId, category)
221+
222+
if (percent < REARM_BELOW) {
223+
if (notified > 0) await rearmThreshold(scope, stateId, category)
224+
return
225+
}
226+
227+
const desired = thresholdFor(percent)
228+
if (params.rearmOnly || desired === 0 || notified >= desired) return
229+
230+
const recipients = await params.resolveRecipients()
231+
if (recipients.length === 0) return
232+
233+
if (!(await claimThreshold(scope, stateId, category, desired))) return
234+
235+
await params.send(desired, recipients)
236+
}
237+
151238
/**
152239
* Send a usage-limit threshold email (80% warning / 100% reached) for a
153-
* non-credit category, edge-triggered on the mutation that changed usage.
240+
* non-credit category, evaluated on the mutation that changed usage.
154241
*
155-
* Flow: bail when billing is off or the limit is non-positive; re-arm the
156-
* persisted threshold when current usage is back in the low band; then (for
157-
* increases) resolve eligible recipients and atomically claim the threshold
158-
* before sending. Re-arm and claim are mutually exclusive per call — re-arm only
159-
* fires when `desired === 0` — so the dedup stays a single atomic
160-
* {@link claimThreshold} with no re-arm/claim interleaving race. Recipients are
161-
* resolved with opt-outs applied BEFORE the claim, so an opted-out recipient
162-
* never burns the threshold (which would suppress a later email once
163-
* notifications are re-enabled). Per-recipient send failures are isolated.
242+
* Bails when billing is off or the limit is non-positive, then runs the shared
243+
* {@link notifyThresholdLevel} flow on the category's persisted threshold.
244+
* Per-recipient send failures are isolated.
164245
*
165246
* The highest threshold already emailed is persisted per category on
166247
* `user_stats` / `organization`; it re-arms once usage drops below
@@ -171,7 +252,7 @@ async function resolveRecipients(
171252
*/
172253
export async function maybeSendLimitThresholdEmail(params: {
173254
category: LimitCategory
174-
scope: 'user' | 'organization'
255+
scope: ThresholdScope
175256
workspaceId: string
176257
currentUsage: number
177258
limit: number
@@ -195,64 +276,66 @@ export async function maybeSendLimitThresholdEmail(params: {
195276
if (params.limit <= 0) return
196277

197278
const { category, scope } = params
198-
const percent = Math.max(0, (params.currentUsage / params.limit) * 100)
199-
const desired = thresholdFor(percent)
200-
201279
const stateId = scope === 'user' ? params.userId : params.organizationId
202280
if (!stateId) return
203281

204-
if (percent < REARM_BELOW) {
205-
await rearmThreshold(scope, stateId, category)
206-
}
282+
const percent = Math.max(0, (params.currentUsage / params.limit) * 100)
207283

208-
if (params.rearmOnly || desired === 0) return
209-
210-
const recipients = await resolveRecipients(scope, params)
211-
if (recipients.length === 0) return
212-
213-
if (!(await claimThreshold(scope, stateId, category, desired))) return
214-
215-
const kind = desired === REACH_THRESHOLD ? 'reached' : 'warning'
216-
const percentUsed = Math.min(100, Math.round(percent))
217-
const upgradeLink = `${getBaseUrl()}${buildUpgradeHref(params.workspaceId, category)}`
218-
219-
const [{ getLimitEmailSubject, renderLimitThresholdEmail }, { sendEmail }] = await Promise.all([
220-
import('@/components/emails'),
221-
import('@/lib/messaging/email/mailer'),
222-
])
223-
224-
let sent = 0
225-
for (const r of recipients) {
226-
try {
227-
const html = await renderLimitThresholdEmail({
228-
kind,
229-
reason: category,
230-
userName: r.name,
231-
usageLabel: params.usageLabel,
232-
limitLabel: params.limitLabel,
233-
percentUsed,
234-
upgradeLink,
235-
})
236-
237-
await sendEmail({
238-
to: r.email,
239-
subject: getLimitEmailSubject(category, kind),
240-
html,
241-
emailType: 'notifications',
242-
})
243-
sent++
244-
} catch (sendError) {
245-
logger.error('Failed to send limit email', {
246-
category,
247-
email: r.email,
248-
error: sendError,
249-
})
250-
}
251-
}
284+
await notifyThresholdLevel({
285+
scope,
286+
stateId,
287+
category,
288+
percent,
289+
rearmOnly: params.rearmOnly,
290+
resolveRecipients: () => resolveLimitEmailRecipients(scope, params),
291+
send: async (threshold, recipients) => {
292+
const kind = threshold === REACH_THRESHOLD ? 'reached' : 'warning'
293+
const percentUsed = Math.min(100, Math.round(percent))
294+
const upgradeLink = `${getBaseUrl()}${buildUpgradeHref(params.workspaceId, category)}`
252295

253-
if (sent > 0) {
254-
logger.info('Sent usage-limit threshold email', { category, scope, kind, percentUsed, sent })
255-
}
296+
const [{ getLimitEmailSubject, renderLimitThresholdEmail }, { sendEmail }] =
297+
await Promise.all([import('@/components/emails'), import('@/lib/messaging/email/mailer')])
298+
299+
let sent = 0
300+
for (const r of recipients) {
301+
try {
302+
const html = await renderLimitThresholdEmail({
303+
kind,
304+
reason: category,
305+
userName: r.name,
306+
usageLabel: params.usageLabel,
307+
limitLabel: params.limitLabel,
308+
percentUsed,
309+
upgradeLink,
310+
})
311+
312+
await sendEmail({
313+
to: r.email,
314+
subject: getLimitEmailSubject(category, kind),
315+
html,
316+
emailType: 'notifications',
317+
})
318+
sent++
319+
} catch (sendError) {
320+
logger.error('Failed to send limit email', {
321+
category,
322+
email: r.email,
323+
error: sendError,
324+
})
325+
}
326+
}
327+
328+
if (sent > 0) {
329+
logger.info('Sent usage-limit threshold email', {
330+
category,
331+
scope,
332+
kind,
333+
percentUsed,
334+
sent,
335+
})
336+
}
337+
},
338+
})
256339
} catch (error) {
257340
logger.error('Failed to send usage-limit threshold email', {
258341
category: params.category,

0 commit comments

Comments
 (0)