Skip to content

Commit cbd2603

Browse files
authored
fix(billing): let a slow ledger read finish instead of blocking execution at the singleflight default (#8146)
* fix(billing): let a slow ledger read finish instead of blocking execution at the singleflight default * fix(billing): bound the ledger sum at the database and derive the gate deadline from it * fix(billing): keep the ledger statement bound in the billing constants module
1 parent 8e9d5fd commit cbd2603

6 files changed

Lines changed: 113 additions & 11 deletions

File tree

‎apps/sim/lib/billing/constants.ts‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,15 @@ export const DEFAULT_OVERAGE_THRESHOLD = 100
3939
*/
4040
export const BILLING_LOCK_TIMEOUT_MS = 5_000
4141

42+
/**
43+
* Bound on one ledger sum. A large payer's period covers millions of rows, and from a cold
44+
* cache or under heavy I/O the sum can run for tens of seconds; past this the database ends it
45+
* and the read fails, so a caller that admits on the sum fails closed rather than waiting
46+
* without limit. The usage gate derives its coalescing deadline from this bound, so the sum
47+
* always ends at the database before the gate gives up on it.
48+
*/
49+
export const USAGE_LEDGER_STATEMENT_TIMEOUT_MS = 60_000
50+
4251
/**
4352
* Available credit tiers. Each tier maps a credit amount to the underlying dollar
4453
* cost and carries that tier's fixed weekly refresh allowance.

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

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import {
1616
checkIngestionUsageLimits,
1717
checkSearchUsageLimits,
1818
resetUsageGateCache,
19+
USAGE_GATE_SETTLE_TIMEOUT_MS,
1920
USAGE_GATE_TTL_MS,
2021
} from '@/lib/billing/core/usage-gate-cache'
2122

@@ -178,4 +179,36 @@ describe('checkExecutionUsageLimits', () => {
178179
await checkExecutionUsageLimits(ATTRIBUTION)
179180
expect(mockCheck).toHaveBeenCalledTimes(2)
180181
})
182+
183+
it('waits for a slow ledger read past the singleflight default instead of blocking', async () => {
184+
vi.useFakeTimers()
185+
try {
186+
mockCheck.mockImplementationOnce(() => sleep(45_000).then(() => ({ isExceeded: false })))
187+
const pending = checkExecutionUsageLimits(ATTRIBUTION)
188+
await vi.advanceTimersByTimeAsync(45_000)
189+
await expect(pending).resolves.toEqual({ isExceeded: false })
190+
expect(mockCheck).toHaveBeenCalledTimes(1)
191+
} finally {
192+
vi.useRealTimers()
193+
}
194+
})
195+
196+
it('gives up on a read that never answers at the gate deadline, then reads fresh', async () => {
197+
vi.useFakeTimers()
198+
try {
199+
mockCheck.mockReturnValueOnce(new Promise(() => {}))
200+
const hung = checkExecutionUsageLimits(ATTRIBUTION)
201+
const rejection = expect(hung).rejects.toThrow(
202+
`did not settle within ${USAGE_GATE_SETTLE_TIMEOUT_MS}ms`
203+
)
204+
await vi.advanceTimersByTimeAsync(USAGE_GATE_SETTLE_TIMEOUT_MS)
205+
await rejection
206+
await expect(checkExecutionUsageLimits(ATTRIBUTION)).resolves.toEqual({
207+
isExceeded: false,
208+
})
209+
expect(mockCheck).toHaveBeenCalledTimes(2)
210+
} finally {
211+
vi.useRealTimers()
212+
}
213+
})
181214
})

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

Lines changed: 19 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { LRUCache } from 'lru-cache'
2+
import { USAGE_LEDGER_STATEMENT_TIMEOUT_MS } from '@/lib/billing/constants'
23
import {
34
type AttributedUsageLimitsResult,
45
type BillingAttributionSnapshot,
@@ -20,6 +21,17 @@ import { coalesceLocally } from '@/lib/concurrency/singleflight'
2021
*/
2122
export const USAGE_GATE_TTL_MS = 5 * 60 * 1000
2223

24+
/**
25+
* How long a coalesced usage read may take before its callers give up on it. The read's cost is
26+
* the ledger sum, which the database ends at {@link USAGE_LEDGER_STATEMENT_TIMEOUT_MS}; the
27+
* remainder is a few indexed lookups and the connection waits around them. The singleflight
28+
* default of 30 s exists to bound a hung producer, and a slow sum is not a hung one: given up on
29+
* early, it keeps running detached while every joined caller fails and the next caller starts a
30+
* second sum alongside it. Derived from the statement bound so the database always ends the sum
31+
* first, and the gate only gives up on a connection that never answers.
32+
*/
33+
export const USAGE_GATE_SETTLE_TIMEOUT_MS = USAGE_LEDGER_STATEMENT_TIMEOUT_MS + 15_000
34+
2335
/**
2436
* Recent gate answers, admitted and refused, with `LRUCache` supplying the TTL
2537
* and the size bound. Each entry point decides which of them it may serve.
@@ -62,9 +74,9 @@ function gateKey(attribution: BillingAttributionSnapshot): string {
6274
* the cache. A read that throws writes nothing.
6375
*
6476
* `coalesceLocally` collapses concurrent misses onto one ledger read and bounds
65-
* a hung read at its settle deadline. The write stays on the value this caller
66-
* received, so a producer that timed out and later resolved cannot overwrite a
67-
* fresher answer.
77+
* a hung read at {@link USAGE_GATE_SETTLE_TIMEOUT_MS}. The write stays on the
78+
* value this caller received, so a producer that timed out and later resolved
79+
* cannot overwrite a fresher answer.
6880
*
6981
* There is deliberately no invalidator: usage and limit changes land in other
7082
* processes (execution workers, Stripe webhooks), so the TTL is the real bound.
@@ -77,8 +89,10 @@ async function checkUsageLimitsThroughCache(
7789
const cached = gateCache.get(key)
7890
if (cached !== undefined && (cacheRefusals || !cached.isExceeded)) return cached
7991

80-
const result = await coalesceLocally(`usage-gate:${key}`, () =>
81-
checkAttributedUsageLimits(attribution)
92+
const result = await coalesceLocally(
93+
`usage-gate:${key}`,
94+
() => checkAttributedUsageLimits(attribution),
95+
USAGE_GATE_SETTLE_TIMEOUT_MS
8296
)
8397
if (cacheRefusals || !result.isExceeded) gateCache.set(key, result)
8498
return result

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

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,9 +33,11 @@ vi.mock('@/lib/billing/subscriptions/utils', () => ({
3333
isOrgScopedSubscription: mockIsOrgScopedSubscription,
3434
}))
3535

36+
import { USAGE_LEDGER_STATEMENT_TIMEOUT_MS } from '@/lib/billing/constants'
3637
import {
3738
CUMULATIVE_COST_EPSILON,
3839
CumulativeUsageContextMismatchError,
40+
getBillingPeriodUsageCost,
3941
getUserUsageLogs,
4042
getWorkspaceUsageLogs,
4143
recordCumulativeUsage,
@@ -554,3 +556,35 @@ describe('usage-log query scopes', () => {
554556
})
555557
})
556558
})
559+
560+
describe('getBillingPeriodUsageCost', () => {
561+
beforeEach(() => {
562+
vi.clearAllMocks()
563+
installSharedDbMocks()
564+
})
565+
566+
it('bounds the ledger sum with its own statement timeout inside one transaction', async () => {
567+
const execute = vi.fn().mockResolvedValue([])
568+
const where = vi.fn().mockResolvedValue([{ cost: '12.5' }])
569+
const tx = { execute, select: vi.fn(() => ({ from: vi.fn(() => ({ where })) })) }
570+
mockTransaction.mockImplementation((callback: (client: typeof tx) => Promise<unknown>) =>
571+
callback(tx)
572+
)
573+
574+
const cost = await getBillingPeriodUsageCost(
575+
{ type: 'organization', id: 'org-1' },
576+
{ start: new Date('2026-05-01T00:00:00Z'), end: new Date('2027-05-01T00:00:00Z') }
577+
)
578+
579+
expect(cost).toBe(12.5)
580+
expect(mockTransaction).toHaveBeenCalledTimes(1)
581+
const executed = execute.mock.calls.map(
582+
([statement]) => (statement as { toSQL: () => { sql: string } }).toSQL().sql
583+
)
584+
expect(executed).toContain(
585+
`SET LOCAL statement_timeout = '${USAGE_LEDGER_STATEMENT_TIMEOUT_MS}ms'`
586+
)
587+
/** The bound is set before the sum runs, not after. */
588+
expect(execute.mock.invocationCallOrder[0]).toBeLessThan(where.mock.invocationCallOrder[0])
589+
})
590+
})

‎apps/sim/lib/billing/core/usage-log.ts‎

Lines changed: 16 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import {
1515
textKey,
1616
timestampKey,
1717
} from '@/lib/api/list-query'
18+
import { USAGE_LEDGER_STATEMENT_TIMEOUT_MS } from '@/lib/billing/constants'
1819
import { defaultBillingPeriod } from '@/lib/billing/core/billing-period'
1920
import { getHighestPrioritySubscription } from '@/lib/billing/core/plan'
2021
import {
@@ -215,6 +216,10 @@ async function resolveBillingContext(
215216
/**
216217
* Returns attributed ledger usage for a billing entity/period. The ledger is
217218
* the sole source of truth for usage — there is no userStats baseline.
219+
*
220+
* The sum runs in a transaction of its own on the given client so that it can
221+
* be bounded by {@link USAGE_LEDGER_STATEMENT_TIMEOUT_MS} for that statement
222+
* alone: `SET LOCAL` ends with the transaction and never reaches the pool.
218223
*/
219224
export async function getBillingPeriodUsageCost(
220225
billingEntity: BillingEntity,
@@ -238,12 +243,17 @@ export async function getBillingPeriodUsageCost(
238243
)
239244
}
240245

241-
const [row] = await executor
242-
.select({
243-
cost: sql<string>`COALESCE(SUM(${usageLog.cost}), 0)`,
244-
})
245-
.from(usageLog)
246-
.where(and(...conditions))
246+
const [row] = await executor.transaction(async (tx) => {
247+
await tx.execute(
248+
sql.raw(`SET LOCAL statement_timeout = '${USAGE_LEDGER_STATEMENT_TIMEOUT_MS}ms'`)
249+
)
250+
return tx
251+
.select({
252+
cost: sql<string>`COALESCE(SUM(${usageLog.cost}), 0)`,
253+
})
254+
.from(usageLog)
255+
.where(and(...conditions))
256+
})
247257

248258
return Number.parseFloat(row?.cost ?? '0')
249259
}

‎apps/sim/lib/billing/enterprise-provisioning.test.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -199,6 +199,8 @@ describe('Enterprise issuance preflight', () => {
199199
queueTableRows(schemaMock.workspace, [])
200200
queueTableRows(schemaMock.workspace, [])
201201
queueTableRows(schemaMock.subscription, [])
202+
/** The run count resolves first; the ledger sum opens its bounded transaction before it reads. */
203+
queueTableRows(schemaMock.usageLog, [{ workflowRuns: 0 }])
202204
queueTableRows(schemaMock.usageLog, [{ cost: '150' }])
203205

204206
const result = await getEnterpriseIssuancePreflight({

0 commit comments

Comments
 (0)