Skip to content

Commit 08c9641

Browse files
committed
fix(billing): cache enterprise reporting-window usage sums for soft usage gates
An enterprise usage period is a reporting window up to a year long, so every soft usage check re-summed the payer's whole year of usage_log, once per billable event (every run admission, Chat request, v1 API response, and knowledge document). For a busy enterprise org that scan ran thousands of times an hour and dominated database CPU. readSoftGateUsageCost serves reporting-window sums from a per-process LRUCache (fetchMethod, max 1000, 30 s TTL) keyed by payer, period source, and window bounds; concurrent misses coalesce and a failed sum is never cached. The raw cached reader is private and typed to reporting periods, so every other period is summed exactly. Pooled org admission (computePooledOrgUsage) and the v1 usage report (getEffectiveCurrentPeriodCost) read through it. The ledger only grows within a window, so a cached sum can only trail the true one by at most one TTL of usage: an admission gate may let a payer run on briefly past its limit, never refuse it wrongly. Exact paths are untouched: getBillingPeriodUsageCost itself (read-your-writes), threshold billing, cycle close, invoicing, analytics, and the execution logger's edge-triggered usage emails, which need an exact baseline.
1 parent df3b413 commit 08c9641

5 files changed

Lines changed: 352 additions & 3 deletions

File tree

‎apps/sim/lib/billing/calculations/usage-monitor.test.ts‎

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,58 @@ describe('checkUsageStatus', () => {
105105
)
106106
})
107107

108+
it('shares one pooled sum across admissions in an enterprise reporting window', async () => {
109+
const billingPeriod = {
110+
start: new Date('2026-01-01T00:00:00.000Z'),
111+
end: new Date('2027-01-01T00:00:00.000Z'),
112+
source: 'reporting' as const,
113+
anchorDate: '2026-01-01',
114+
interval: 'year' as const,
115+
}
116+
const subscription = {
117+
referenceId: 'org-reporting-shared',
118+
plan: 'enterprise',
119+
status: 'active',
120+
seats: 1,
121+
periodStart: billingPeriod.start,
122+
periodEnd: billingPeriod.end,
123+
}
124+
const billingContext = {
125+
billingEntity: { type: 'organization' as const, id: 'org-reporting-shared' },
126+
billingPeriod,
127+
}
128+
129+
await checkUsageStatus('user-1', subscription, billingContext)
130+
await checkUsageStatus('user-2', subscription, billingContext)
131+
132+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(1)
133+
})
134+
135+
it('sums a Stripe-period organization pool exactly on every admission', async () => {
136+
const billingPeriod = {
137+
start: new Date('2026-06-01T00:00:00.000Z'),
138+
end: new Date('2026-07-01T00:00:00.000Z'),
139+
source: 'stripe' as const,
140+
}
141+
const subscription = {
142+
referenceId: 'org-stripe',
143+
plan: 'team',
144+
status: 'active',
145+
seats: 1,
146+
periodStart: null,
147+
periodEnd: null,
148+
}
149+
const billingContext = {
150+
billingEntity: { type: 'organization' as const, id: 'org-stripe' },
151+
billingPeriod,
152+
}
153+
154+
await checkUsageStatus('user-1', subscription, billingContext)
155+
await checkUsageStatus('user-1', subscription, billingContext)
156+
157+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(2)
158+
})
159+
108160
it('reads paid personal ledger usage and refresh from one snapshot', async () => {
109161
const periodStart = new Date('2026-06-01T00:00:00.000Z')
110162
const periodEnd = new Date('2026-07-01T00:00:00.000Z')

‎apps/sim/lib/billing/calculations/usage-monitor.ts‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import { isOrganizationBillingBlocked } from '@/lib/billing/core/access'
77
import { defaultBillingPeriod } from '@/lib/billing/core/billing-period'
88
import { getHighestPrioritySubscription } from '@/lib/billing/core/plan'
99
import { resolveSubscriptionUsagePeriod } from '@/lib/billing/core/reporting-period'
10+
import { readSoftGateUsageCost } from '@/lib/billing/core/reporting-usage-cache'
1011
import { getUserUsageLimit, type UsageLimitSubscription } from '@/lib/billing/core/usage'
1112
import {
1213
type BillingContext,
@@ -44,6 +45,11 @@ interface UsageData {
4445
organizationId: string | null
4546
}
4647

48+
/**
49+
* The organization's pooled usage for an admission check. An enterprise reporting window is
50+
* served through {@link readSoftGateUsageCost}; its weekly refresh is zero, so it always takes
51+
* a plain-sum branch below.
52+
*/
4753
async function computePooledOrgUsage(
4854
organizationId: string,
4955
sub: UsageLimitSubscription,
@@ -58,12 +64,12 @@ async function computePooledOrgUsage(
5864
}
5965

6066
if (!isPaid(sub.plan) || !sub.periodStart) {
61-
return getBillingPeriodUsageCost({ type: 'organization', id: organizationId }, billingPeriod)
67+
return readSoftGateUsageCost({ type: 'organization', id: organizationId }, billingPeriod)
6268
}
6369

6470
const weeklyRefreshDollars = getPlanWeeklyRefreshDollars(sub.plan)
6571
if (weeklyRefreshDollars <= 0) {
66-
return getBillingPeriodUsageCost({ type: 'organization', id: organizationId }, billingPeriod)
72+
return readSoftGateUsageCost({ type: 'organization', id: organizationId }, billingPeriod)
6773
}
6874

6975
const { ledgerUsage, refreshConsumed } = await computeBillingPeriodUsageWithWeeklyRefresh({
Lines changed: 197 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,197 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { db } from '@sim/db'
5+
import { sleep } from '@sim/utils/helpers'
6+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
7+
8+
const { mockGetBillingPeriodUsageCost } = vi.hoisted(() => ({
9+
mockGetBillingPeriodUsageCost: vi.fn(),
10+
}))
11+
12+
vi.mock('@/lib/billing/core/usage-log', () => ({
13+
getBillingPeriodUsageCost: mockGetBillingPeriodUsageCost,
14+
}))
15+
16+
import * as reportingUsageCache from '@/lib/billing/core/reporting-usage-cache'
17+
import type { UsageQueryPeriod } from '@/lib/billing/core/usage-log'
18+
19+
const { REPORTING_USAGE_CACHE_TTL_MS, readSoftGateUsageCost } = reportingUsageCache
20+
21+
const REPORTING: UsageQueryPeriod = {
22+
start: new Date('2026-01-01T00:00:00.000Z'),
23+
end: new Date('2027-01-01T00:00:00.000Z'),
24+
source: 'reporting',
25+
}
26+
const NEXT_REPORTING: UsageQueryPeriod = {
27+
start: new Date('2027-01-01T00:00:00.000Z'),
28+
end: new Date('2028-01-01T00:00:00.000Z'),
29+
source: 'reporting',
30+
}
31+
const STRIPE: UsageQueryPeriod = {
32+
start: new Date('2026-09-01T00:00:00.000Z'),
33+
end: new Date('2026-10-01T00:00:00.000Z'),
34+
source: 'stripe',
35+
}
36+
37+
let nextOrg = 0
38+
/** A fresh payer per test, since the cache is module state shared across tests. */
39+
function freshOrg() {
40+
nextOrg += 1
41+
return { type: 'organization' as const, id: `org-${nextOrg}` }
42+
}
43+
44+
describe('readSoftGateUsageCost on a reporting window', () => {
45+
beforeEach(() => {
46+
vi.clearAllMocks()
47+
mockGetBillingPeriodUsageCost.mockReset()
48+
})
49+
50+
afterEach(() => {
51+
vi.restoreAllMocks()
52+
})
53+
54+
it('sums the ledger again once a cached sum outlives its TTL', async () => {
55+
const org = freshOrg()
56+
const start = performance.now()
57+
const clock = vi.spyOn(performance, 'now').mockReturnValue(start)
58+
mockGetBillingPeriodUsageCost.mockResolvedValueOnce(10).mockResolvedValueOnce(25)
59+
60+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(10)
61+
/** The cache re-reads its clock only after `ttlResolution` (1 ms) of real time. */
62+
clock.mockReturnValue(start + REPORTING_USAGE_CACHE_TTL_MS - 1)
63+
await sleep(2)
64+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(10)
65+
clock.mockReturnValue(start + REPORTING_USAGE_CACHE_TTL_MS + 1)
66+
await sleep(2)
67+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(25)
68+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(2)
69+
})
70+
71+
it('coalesces concurrent and repeated reads of one window into one ledger sum', async () => {
72+
const org = freshOrg()
73+
let resolveSum: (value: number) => void = () => {}
74+
mockGetBillingPeriodUsageCost.mockReturnValueOnce(
75+
new Promise<number>((resolve) => {
76+
resolveSum = resolve
77+
})
78+
)
79+
80+
const concurrent = Promise.all([
81+
readSoftGateUsageCost(org, REPORTING),
82+
readSoftGateUsageCost(org, REPORTING),
83+
readSoftGateUsageCost(org, REPORTING),
84+
])
85+
resolveSum(42)
86+
87+
await expect(concurrent).resolves.toEqual([42, 42, 42])
88+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(42)
89+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(1)
90+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledWith(org, REPORTING)
91+
})
92+
93+
it('serves a zero sum from cache rather than re-reading it', async () => {
94+
const org = freshOrg()
95+
mockGetBillingPeriodUsageCost.mockResolvedValue(0)
96+
97+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(0)
98+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(0)
99+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(1)
100+
})
101+
102+
it('keeps separate sums for different payers and windows', async () => {
103+
const first = freshOrg()
104+
const second = freshOrg()
105+
mockGetBillingPeriodUsageCost
106+
.mockResolvedValueOnce(1)
107+
.mockResolvedValueOnce(2)
108+
.mockResolvedValueOnce(3)
109+
.mockResolvedValueOnce(4)
110+
111+
await expect(readSoftGateUsageCost(first, REPORTING)).resolves.toBe(1)
112+
await expect(readSoftGateUsageCost(second, REPORTING)).resolves.toBe(2)
113+
await expect(readSoftGateUsageCost(first, NEXT_REPORTING)).resolves.toBe(3)
114+
await expect(readSoftGateUsageCost({ type: 'user', id: first.id }, REPORTING)).resolves.toBe(4)
115+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(4)
116+
})
117+
118+
it('surfaces a failed sum to every waiting caller and never caches it', async () => {
119+
const org = freshOrg()
120+
const failure = new Error('canceling statement due to statement timeout')
121+
mockGetBillingPeriodUsageCost.mockRejectedValueOnce(failure).mockResolvedValueOnce(17)
122+
123+
const results = await Promise.allSettled([
124+
readSoftGateUsageCost(org, REPORTING),
125+
readSoftGateUsageCost(org, REPORTING),
126+
])
127+
expect(results).toEqual([
128+
{ status: 'rejected', reason: failure },
129+
{ status: 'rejected', reason: failure },
130+
])
131+
132+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(17)
133+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(2)
134+
})
135+
})
136+
137+
describe('readSoftGateUsageCost', () => {
138+
beforeEach(() => {
139+
vi.clearAllMocks()
140+
mockGetBillingPeriodUsageCost.mockReset()
141+
})
142+
143+
it('exposes no cached reader that a non-reporting period could reach', () => {
144+
expect(Object.keys(reportingUsageCache).sort()).toEqual([
145+
'REPORTING_USAGE_CACHE_TTL_MS',
146+
'readSoftGateUsageCost',
147+
])
148+
})
149+
150+
it('never serves a cached reporting sum to another source with the same bounds', async () => {
151+
const org = freshOrg()
152+
const sameBounds = { start: REPORTING.start, end: REPORTING.end }
153+
mockGetBillingPeriodUsageCost.mockResolvedValueOnce(10).mockResolvedValueOnce(99)
154+
155+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(10)
156+
await expect(readSoftGateUsageCost(org, { ...sameBounds, source: 'stripe' })).resolves.toBe(99)
157+
await expect(readSoftGateUsageCost(org, REPORTING)).resolves.toBe(10)
158+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(2)
159+
})
160+
161+
it('serves reporting windows from the cache', async () => {
162+
const org = freshOrg()
163+
mockGetBillingPeriodUsageCost.mockResolvedValue(10)
164+
165+
await readSoftGateUsageCost(org, REPORTING)
166+
await readSoftGateUsageCost(org, REPORTING)
167+
168+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(1)
169+
})
170+
171+
it.each([
172+
['stripe', STRIPE],
173+
['default', { ...STRIPE, source: 'default' as const }],
174+
['unlabelled', { start: STRIPE.start, end: STRIPE.end }],
175+
])('sums %s periods exactly on every call', async (_label, period) => {
176+
const org = freshOrg()
177+
mockGetBillingPeriodUsageCost.mockResolvedValue(10)
178+
179+
await readSoftGateUsageCost(org, period)
180+
await readSoftGateUsageCost(org, period)
181+
182+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(2)
183+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledWith(org, period, undefined, db)
184+
})
185+
186+
it('reads a reporting window exactly on a caller-supplied executor', async () => {
187+
const org = freshOrg()
188+
const executor = { transaction: vi.fn() } as unknown as typeof db
189+
mockGetBillingPeriodUsageCost.mockResolvedValue(10)
190+
191+
await readSoftGateUsageCost(org, REPORTING, executor)
192+
await readSoftGateUsageCost(org, REPORTING, executor)
193+
194+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledTimes(2)
195+
expect(mockGetBillingPeriodUsageCost).toHaveBeenCalledWith(org, REPORTING, undefined, executor)
196+
})
197+
})
Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
import { db } from '@sim/db'
2+
import { LRUCache } from 'lru-cache'
3+
import {
4+
type BillingEntity,
5+
getBillingPeriodUsageCost,
6+
type UsageQueryPeriod,
7+
} from '@/lib/billing/core/usage-log'
8+
import type { DbClient } from '@/lib/db/types'
9+
10+
/**
11+
* How long a reporting-window usage sum is served before it is summed again.
12+
*
13+
* A reporting window is an enterprise contract period, up to a year long, so its sum scans every
14+
* ledger row the payer wrote in that year, and the soft gates below re-ran it once per billable
15+
* event. The ledger only grows within a window (rows are inserted at a cost above zero, and the
16+
* one update is a monotonic top-up), so a served sum is never above the true one: it omits at
17+
* most the usage written since it was read. That is the safe direction for every reader here —
18+
* an admission gate lets a payer run on for at most this long past their limit, and a sum at or
19+
* above the limit is a refusal the true sum would also give. Thirty seconds keeps that overrun
20+
* small against a year-long allowance while turning a per-event scan into one per window.
21+
*/
22+
export const REPORTING_USAGE_CACHE_TTL_MS = 30_000
23+
24+
/** A usage window known to be an enterprise reporting window — the only kind this cache serves. */
25+
type ReportingQueryPeriod = UsageQueryPeriod & { source: 'reporting' }
26+
27+
/**
28+
* Sums shared across callers, one per payer and window. Every key is an enterprise payer's
29+
* current window, a few dozen bytes each, so the ceiling sits far above any process's working
30+
* set and only backstops memory; an eviction inside the TTL costs one extra sum.
31+
*
32+
* `fetchMethod` coalesces concurrent misses onto one sum. A rejected sum is evicted rather than
33+
* stored (`noDeleteOnFetchRejection` and `allowStaleOnFetchRejection` stay off), so every caller
34+
* of that read sees the error it would have seen uncached and the next call sums again. There is
35+
* no settle deadline: the sum runs under the ledger's own `statement_timeout`, so the database
36+
* ends a slow one. There is deliberately no invalidator either — usage is written by execution
37+
* workers in other processes, so the TTL is the real bound.
38+
*/
39+
const reportingUsageCache = new LRUCache<
40+
string,
41+
number,
42+
{ entity: BillingEntity; period: ReportingQueryPeriod }
43+
>({
44+
max: 1_000,
45+
ttl: REPORTING_USAGE_CACHE_TTL_MS,
46+
fetchMethod: (_key, _stale, { context }) =>
47+
getBillingPeriodUsageCost(context.entity, context.period),
48+
})
49+
50+
/**
51+
* The key names the period's source as well as its bounds, so a sum can only ever be shared
52+
* with a read of the same kind of window, even if another source someday reaches this cache.
53+
*/
54+
function reportingUsageKey(entity: BillingEntity, period: ReportingQueryPeriod): string {
55+
return `${entity.type}:${entity.id}:${period.source}:${period.start.toISOString()}:${period.end.toISOString()}`
56+
}
57+
58+
function isReportingPeriod(period: UsageQueryPeriod): period is ReportingQueryPeriod {
59+
return period.source === 'reporting'
60+
}
61+
62+
async function readCachedReportingUsageCost(
63+
entity: BillingEntity,
64+
period: ReportingQueryPeriod
65+
): Promise<number> {
66+
const cost = await reportingUsageCache.fetch(reportingUsageKey(entity, period), {
67+
context: { entity, period },
68+
})
69+
return cost !== undefined ? cost : getBillingPeriodUsageCost(entity, period)
70+
}
71+
72+
/**
73+
* Period usage for a soft reader: an admission check, a display, or a level-triggered
74+
* notification that tolerates the cache's bounded under-count. Enterprise reporting windows are
75+
* served from the shared cache for up to {@link REPORTING_USAGE_CACHE_TTL_MS}, since their
76+
* year-long sum is the expensive one; every other period is summed exactly, as before. A read on
77+
* a caller's own executor (a transaction or a replica) keeps its own snapshot and is never shared.
78+
* Never use it for invoicing, cycle close, an edge-triggered decision, or a read that must see its
79+
* own write.
80+
*/
81+
export function readSoftGateUsageCost(
82+
entity: BillingEntity,
83+
period: UsageQueryPeriod,
84+
executor: DbClient = db
85+
): Promise<number> {
86+
if (isReportingPeriod(period) && executor === db) {
87+
return readCachedReportingUsageCost(entity, period)
88+
}
89+
return getBillingPeriodUsageCost(entity, period, undefined, executor)
90+
}

0 commit comments

Comments
 (0)