Skip to content

Commit 655432c

Browse files
committed
improvement(billing): key the usage email claim on period and limit, and harden the shared read
1 parent ea79a8b commit 655432c

8 files changed

Lines changed: 203 additions & 137 deletions

File tree

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

Lines changed: 54 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +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 { and, eq, sql } from 'drizzle-orm'
5+
import { and, eq, type SQL, sql } from 'drizzle-orm'
66
import type { HighestPrioritySubscription } from '@/lib/billing/core/plan'
77
import { getHighestPrioritySubscription } from '@/lib/billing/core/subscription'
88
import type { BillingEntity } from '@/lib/billing/core/usage-log'
@@ -17,9 +17,6 @@ const logger = createLogger('LimitNotifications')
1717
/** Limit categories that send per-category threshold emails (credits has its own path). */
1818
export type LimitCategory = Extract<UpgradeReason, 'storage' | 'tables' | 'seats'>
1919

20-
/** Every category whose emailed threshold is persisted, including credits. */
21-
type ClaimCategory = LimitCategory | Extract<UpgradeReason, 'credits'>
22-
2320
const WARN_THRESHOLD = 80
2421
const REACH_THRESHOLD = 100
2522
/** Usage must drop below this band before the same threshold can re-notify (hysteresis). */
@@ -35,6 +32,32 @@ function thresholdFor(percent: number): 0 | 80 | 100 {
3532
return 0
3633
}
3734

35+
/**
36+
* Replace the account's `limitNotifications` with `next` when `condition` holds, returning
37+
* whether the row was updated.
38+
*/
39+
async function writeLimitNotifications(
40+
scope: 'user' | 'organization',
41+
id: string,
42+
next: SQL,
43+
condition: SQL
44+
): Promise<boolean> {
45+
const written =
46+
scope === 'user'
47+
? await db
48+
.update(userStats)
49+
.set({ limitNotifications: next })
50+
.where(and(eq(userStats.userId, id), condition))
51+
.returning({ id: userStats.userId })
52+
: await db
53+
.update(organization)
54+
.set({ limitNotifications: next })
55+
.where(and(eq(organization.id, id), condition))
56+
.returning({ id: organization.id })
57+
58+
return written.length > 0
59+
}
60+
3861
/**
3962
* Atomically claim a threshold for a category: advance the stored value to
4063
* `threshold` only if it is currently lower, returning whether THIS call won the
@@ -44,7 +67,7 @@ function thresholdFor(percent: number): 0 | 80 | 100 {
4467
async function claimThreshold(
4568
scope: 'user' | 'organization',
4669
id: string,
47-
category: ClaimCategory,
70+
category: LimitCategory,
4871
threshold: number
4972
): Promise<boolean> {
5073
const setExpr = sql`jsonb_set(coalesce(${scope === 'user' ? userStats.limitNotifications : organization.limitNotifications}, '{}'::jsonb), ARRAY[${category}], to_jsonb(${threshold}::int))`
@@ -53,38 +76,37 @@ async function claimThreshold(
5376
? sql`coalesce((${userStats.limitNotifications} ->> ${category})::int, 0) < ${threshold}`
5477
: sql`coalesce((${organization.limitNotifications} ->> ${category})::int, 0) < ${threshold}`
5578

56-
const claimed =
57-
scope === 'user'
58-
? await db
59-
.update(userStats)
60-
.set({ limitNotifications: setExpr })
61-
.where(and(eq(userStats.userId, id), onlyIfLower))
62-
.returning({ id: userStats.userId })
63-
: await db
64-
.update(organization)
65-
.set({ limitNotifications: setExpr })
66-
.where(and(eq(organization.id, id), onlyIfLower))
67-
.returning({ id: organization.id })
68-
69-
return claimed.length > 0
79+
return writeLimitNotifications(scope, id, setExpr, onlyIfLower)
7080
}
7181

7282
const DAY_MS = 24 * 60 * 60 * 1000
7383

7484
/**
75-
* Claim a credits threshold (80 or 100) once per billing period, returning whether THIS call won
76-
* it. The stored value is the period's start day followed by the threshold, so it only grows: a
77-
* later period outranks every claim of an earlier one and re-arms both thresholds with no reset
78-
* write, while within a period a claim of 100 also retires 80, and never the reverse.
85+
* Claim a credits threshold (80 or 100), returning whether THIS call won it. The claim is keyed
86+
* on the billing period and the limit: `credits` holds the highest threshold emailed while
87+
* `creditsPeriod` (the period's start day) and `creditsLimit` (the limit in cents) still match,
88+
* so a new period or a changed limit — in either direction — re-arms both thresholds with no
89+
* reset write, while within one a claim of 100 also retires 80, and never the reverse.
7990
*/
80-
export function claimCreditsThreshold(
81-
scope: 'user' | 'organization',
82-
id: string,
83-
periodStart: Date,
91+
export function claimCreditsThreshold(params: {
92+
scope: 'user' | 'organization'
93+
id: string
94+
periodStart: Date
95+
limit: number
8496
threshold: 80 | 100
85-
): Promise<boolean> {
86-
const periodDay = Math.floor(periodStart.getTime() / DAY_MS)
87-
return claimThreshold(scope, id, 'credits', periodDay * 1000 + threshold)
97+
}): Promise<boolean> {
98+
const { scope, id, threshold } = params
99+
const periodDay = Math.floor(params.periodStart.getTime() / DAY_MS)
100+
const limitCents = Math.round(params.limit * 100)
101+
const column = scope === 'user' ? userStats.limitNotifications : organization.limitNotifications
102+
const next = sql`coalesce(${column}, '{}'::jsonb) || jsonb_build_object('credits', ${threshold}::int, 'creditsPeriod', ${periodDay}::bigint, 'creditsLimit', ${limitCents}::bigint)`
103+
const unclaimed = sql`not (
104+
(${column} ->> 'creditsPeriod')::bigint is not distinct from ${periodDay}::bigint
105+
and (${column} ->> 'creditsLimit')::bigint is not distinct from ${limitCents}::bigint
106+
and coalesce((${column} ->> 'credits')::int, 0) >= ${threshold}::int
107+
)`
108+
109+
return writeLimitNotifications(scope, id, next, unclaimed)
88110
}
89111

90112
/** Re-arm a category (reset its stored threshold to 0) once usage falls back into the low band. */
@@ -129,7 +151,7 @@ async function isUnsubscribed(email: string): Promise<boolean> {
129151
* Returning an empty list means "nobody to notify" — the caller then skips the
130152
* claim so the dedup state isn't burned without an email going out.
131153
*/
132-
async function resolveRecipients(
154+
export async function resolveLimitEmailRecipients(
133155
scope: 'user' | 'organization',
134156
params: { userId?: string; userEmail?: string; userName?: string; organizationId?: string }
135157
): Promise<LimitEmailRecipient[]> {
@@ -228,7 +250,7 @@ export async function maybeSendLimitThresholdEmail(params: {
228250

229251
if (params.rearmOnly || desired === 0) return
230252

231-
const recipients = await resolveRecipients(scope, params)
253+
const recipients = await resolveLimitEmailRecipients(scope, params)
232254
if (recipients.length === 0) return
233255

234256
if (!(await claimThreshold(scope, stateId, category, desired))) return

‎apps/sim/lib/billing/core/reporting-usage-cache.integration.ts‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -95,11 +95,30 @@ describe.runIf(Boolean(redisUrl))('shared reporting usage read', () => {
9595
expect(ttl).toBeLessThanOrEqual(35_000)
9696
})
9797

98-
it('treats an unreadable stored sum as a miss and overwrites it', async () => {
98+
it('treats an unreadable stored sum as a miss', async () => {
9999
await redis.set(sharedKey(payer), 'not-a-number', 'PX', 30_000)
100100

101101
await expect(readSoftGateUsageCost(payer, REPORTING)).resolves.toBe(5.75)
102-
await vi.waitFor(async () => expect(await redis.get(sharedKey(payer))).toBe('5.75'))
102+
expect(transaction).toHaveBeenCalledTimes(1)
103+
})
104+
105+
it('never lets a slower, older sum replace one stored while it ran', async () => {
106+
transaction.mockImplementationOnce(async (callback) => {
107+
const result = await database.transaction(callback)
108+
await redis.set(sharedKey(payer), '9.99', 'PX', 30_000)
109+
return result
110+
})
111+
112+
await expect(readSoftGateUsageCost(payer, REPORTING)).resolves.toBe(5.75)
113+
expect(await redis.get(sharedKey(payer))).toBe('9.99')
114+
})
115+
116+
it('sums the ledger when the Redis client cannot be built', async () => {
117+
redisConfigMockFns.mockGetRedisClient.mockImplementation(() => {
118+
throw new Error('Invalid Redis configuration')
119+
})
120+
121+
await expect(readSoftGateUsageCost(payer, REPORTING)).resolves.toBe(5.75)
103122
})
104123

105124
it('sums the ledger promptly when Redis is unreachable', async () => {

‎apps/sim/lib/billing/core/reporting-usage-cache.ts‎

Lines changed: 29 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -26,8 +26,10 @@ const logger = createLogger('ReportingUsageCache')
2626
* above the limit is a refusal the true sum would also give. Thirty seconds keeps that overrun
2727
* small against a year-long allowance while turning a per-event scan into one per window.
2828
*
29-
* A sum is held both in Redis, shared by every process, and in each process that reads it, so a
30-
* served sum can be up to twice this old (plus the Redis expiry's jitter).
29+
* A sum is held both in Redis, shared by every process, and in each process that reads it. It
30+
* reflects the ledger as of the moment its sum began, so a served sum can omit usage written over
31+
* the sum's own duration, plus up to this TTL and its jitter in Redis, plus up to this TTL again
32+
* in the reading process.
3133
*/
3234
export const REPORTING_USAGE_CACHE_TTL_MS = 30_000
3335

@@ -57,9 +59,9 @@ function sharedReportingUsageKey(key: string): string {
5759
* sums the ledger, so the cache can cost a read its latency but never its answer.
5860
*/
5961
async function readSharedReportingUsageCost(key: string): Promise<number | undefined> {
60-
const redis = getRedisClient()
61-
if (!redis) return undefined
6262
try {
63+
const redis = getRedisClient()
64+
if (!redis) return undefined
6365
const stored = await withinDeadline(
6466
() => redis.get(sharedReportingUsageKey(key)),
6567
Date.now() + SHARED_READ_TIMEOUT_MS
@@ -76,14 +78,26 @@ async function readSharedReportingUsageCost(key: string): Promise<number | undef
7678
return undefined
7779
}
7880

79-
/** Fire-and-forget: a read never waits on, or fails because of, the shared write. */
81+
function warnSharedWriteFailed(error: unknown): void {
82+
logger.warn('Shared reporting usage write failed', { error: getErrorMessage(error) })
83+
}
84+
85+
/**
86+
* Fire-and-forget: a read never waits on, or fails because of, the shared write. The write only
87+
* lands when no sum is stored (`NX`), so a slow, older sum can never replace a fresher one or
88+
* extend its expiry.
89+
*/
8090
function writeSharedReportingUsageCost(key: string, cost: number): void {
81-
const redis = getRedisClient()
82-
if (!redis) return
83-
const ttlMs = REPORTING_USAGE_CACHE_TTL_MS + randomInt(0, SHARED_TTL_JITTER_MS)
84-
redis.set(sharedReportingUsageKey(key), String(cost), 'PX', ttlMs).catch((error: unknown) => {
85-
logger.warn('Shared reporting usage write failed', { error: getErrorMessage(error) })
86-
})
91+
try {
92+
const redis = getRedisClient()
93+
if (!redis) return
94+
const ttlMs = REPORTING_USAGE_CACHE_TTL_MS + randomInt(0, SHARED_TTL_JITTER_MS)
95+
redis
96+
.set(sharedReportingUsageKey(key), String(cost), 'PX', ttlMs, 'NX')
97+
.catch(warnSharedWriteFailed)
98+
} catch (error) {
99+
warnSharedWriteFailed(error)
100+
}
87101
}
88102

89103
/**
@@ -153,9 +167,10 @@ async function readCachedReportingUsageCost(
153167
/**
154168
* Period usage for a soft reader: an admission check, a display, or a level-triggered
155169
* notification that tolerates the cache's bounded under-count. Enterprise reporting windows are
156-
* served from the shared cache for up to twice {@link REPORTING_USAGE_CACHE_TTL_MS}, since their
157-
* year-long sum is the expensive one; every other period is summed exactly, as before. A read on
158-
* a caller's own executor (a transaction or a replica) keeps its own snapshot and is never shared.
170+
* served from the shared cache, within the lag {@link REPORTING_USAGE_CACHE_TTL_MS} describes,
171+
* since their year-long sum is the expensive one; every other period is summed exactly, as
172+
* before. A read on a caller's own executor (a transaction or a replica) keeps its own snapshot
173+
* and is never shared.
159174
* Never use it for invoicing, cycle close, an edge-triggered decision, or a read that must see its
160175
* own write.
161176
*/

‎apps/sim/lib/billing/core/usage-threshold-email.integration.ts‎

Lines changed: 46 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -40,27 +40,37 @@ const SEPTEMBER = new Date('2026-09-01T00:00:00.000Z')
4040
const OCTOBER = new Date('2026-10-01T00:00:00.000Z')
4141

4242
let organizationId: string
43+
let adminId: string
4344

44-
function notify(currentUsage: number, periodStart = SEPTEMBER) {
45+
/** One completion that leaves the organization at `usage`, having recorded `costDelta` of it. */
46+
function notify(
47+
usage: number,
48+
{ periodStart = SEPTEMBER, limit = 100, costDelta = 1 } = {}
49+
): Promise<void> {
4550
return maybeSendUsageThresholdEmail({
4651
scope: 'organization',
4752
organizationId,
4853
planName: 'Enterprise',
4954
periodStart,
5055
workspaceId: 'workspace',
51-
currentUsage,
52-
limit: 100,
56+
usageBefore: usage - costDelta,
57+
costDelta,
58+
limit,
5359
})
5460
}
5561

62+
async function setNotificationsEnabled(enabled: boolean): Promise<void> {
63+
await connection`INSERT INTO settings VALUES (${adminId}, ${adminId}, ${enabled})
64+
ON CONFLICT (id) DO UPDATE SET billing_usage_notifications_enabled = ${enabled}`
65+
}
66+
5667
beforeAll(async () => {
5768
await connection.unsafe(`CREATE SCHEMA "${schemaName}"`)
5869
await connection.unsafe(`
5970
CREATE TABLE organization (id text PRIMARY KEY, limit_notifications jsonb);
6071
CREATE TABLE "user" (id text PRIMARY KEY, email text, name text);
6172
CREATE TABLE member (id text PRIMARY KEY, organization_id text, user_id text, role text);
6273
CREATE TABLE settings (id text PRIMARY KEY, user_id text, billing_usage_notifications_enabled boolean);
63-
INSERT INTO "user" VALUES ('admin', 'admin@example.com', 'Admin');
6474
`)
6575
select.mockImplementation((fields) => database.select(fields))
6676
update.mockImplementation((table) => database.update(table))
@@ -70,8 +80,10 @@ beforeAll(async () => {
7080
beforeEach(async () => {
7181
mockSendEmail.mockClear()
7282
organizationId = generateId()
83+
adminId = generateId()
7384
await connection`INSERT INTO organization (id) VALUES (${organizationId})`
74-
await connection`INSERT INTO member VALUES (${generateId()}, ${organizationId}, 'admin', 'owner')`
85+
await connection`INSERT INTO "user" VALUES (${adminId}, ${`${adminId}@example.com`}, 'Admin')`
86+
await connection`INSERT INTO member VALUES (${generateId()}, ${organizationId}, ${adminId}, 'owner')`
7587
})
7688

7789
afterAll(async () => {
@@ -97,12 +109,37 @@ describe('usage threshold email', () => {
97109
expect(mockSendEmail).toHaveBeenCalledTimes(2)
98110
})
99111

100-
it('re-arms both thresholds in the next billing period', async () => {
112+
it('re-arms both thresholds whenever the billing period changes, even to an earlier one', async () => {
101113
await notify(100)
102-
await notify(85, OCTOBER)
103-
await notify(100, OCTOBER)
114+
await notify(85, { periodStart: OCTOBER })
115+
await notify(100, { periodStart: OCTOBER })
104116
await notify(85)
105117

106-
expect(mockSendEmail).toHaveBeenCalledTimes(3)
118+
expect(mockSendEmail).toHaveBeenCalledTimes(4)
119+
})
120+
121+
it('warns again at a raised limit after the old one was reached', async () => {
122+
await notify(100)
123+
await notify(100, { limit: 125 })
124+
await notify(110, { limit: 125 })
125+
126+
expect(mockSendEmail).toHaveBeenCalledTimes(2)
127+
})
128+
129+
it('keeps the claim for a later completion when nobody can be notified', async () => {
130+
await setNotificationsEnabled(false)
131+
await notify(90)
132+
await setNotificationsEnabled(true)
133+
await notify(90)
134+
135+
expect(mockSendEmail).toHaveBeenCalledTimes(1)
136+
})
137+
138+
it('keeps the claim when a completion recorded no cost', async () => {
139+
await notify(90, { costDelta: 0 })
140+
expect(mockSendEmail).not.toHaveBeenCalled()
141+
142+
await notify(90)
143+
expect(mockSendEmail).toHaveBeenCalledTimes(1)
107144
})
108145
})

‎apps/sim/lib/billing/core/usage.test.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -334,7 +334,8 @@ describe('maybeSendUsageThresholdEmail', () => {
334334
it('emails a paid personal account at 100% with the raise-your-limit template', async () => {
335335
await maybeSendUsageThresholdEmail({
336336
...paidUser,
337-
currentUsage: 20,
337+
usageBefore: 19,
338+
costDelta: 1,
338339
})
339340

340341
expect(mockRenderUsageLimitReached).toHaveBeenCalledWith(
@@ -359,7 +360,8 @@ describe('maybeSendUsageThresholdEmail', () => {
359360
organizationId: 'org-1',
360361
workspaceId: 'ws-1',
361362
periodStart: new Date('2026-09-01T00:00:00.000Z'),
362-
currentUsage: 500,
363+
usageBefore: 499,
364+
costDelta: 1,
363365
limit: 500,
364366
})
365367

0 commit comments

Comments
 (0)