Skip to content

Commit e1d1102

Browse files
committed
fix(billing): cache admitted direct-v1 continuation verdicts and bound the callback's standing read
- A direct-v1 continuation re-read the payer's full period ledger on every resume leg. Its admitted verdict is now served for the execution gate's TTL, with concurrent misses coalesced, like the attributed path; a refusal or an unreadable ledger is always read again, and the read is skipped when billing is off. - The cost callback waits at most 1 s on the payer's standing, well inside the worker's 5 s callback timeout, and answers not exceeded past it; the abandoned read still caches its admission. - The straddling-period test now reaches the re-judge branch.
1 parent 3a7134d commit e1d1102

5 files changed

Lines changed: 118 additions & 20 deletions

File tree

‎apps/sim/app/api/billing/update-cost/route.test.ts‎

Lines changed: 20 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ vi.mock('@/lib/billing/threshold-billing', () => ({
4949
}))
5050

5151
import { billingUpdateCostResponseSchema } from '@/lib/api/contracts/subscription'
52-
import { resetMidRunPeriodCache } from '@/lib/billing/core/mid-run-usage'
52+
import { resetMidRunUsageCaches } from '@/lib/billing/core/mid-run-usage'
5353
import { resetUsageGateCache } from '@/lib/billing/core/usage-gate-cache'
5454
import {
5555
BillingCallbackBody,
@@ -897,7 +897,7 @@ describe('POST /api/billing/update-cost — mid-run usage gate', () => {
897897

898898
beforeEach(() => {
899899
resetUsageGateCache()
900-
resetMidRunPeriodCache()
900+
resetMidRunUsageCaches()
901901
setEnvFlags({ isBillingEnabled: true, isHosted: true })
902902
mockCheckInternalApiKey.mockReturnValue({ success: true })
903903
mockRecordCumulativeUsage.mockResolvedValue({ billed: true, delta: 0.5, total: 0.5 })
@@ -1136,8 +1136,9 @@ describe('POST /api/billing/update-cost — mid-run usage gate', () => {
11361136
end: new Date(Date.now() + 40).toISOString(),
11371137
},
11381138
}
1139-
mockRequireBillingAttributionHeader.mockReturnValue(straddling)
1140-
mockRefreshAttributionPeriod.mockResolvedValue(CURRENT_ATTRIBUTION)
1139+
mockRefreshAttributionPeriod
1140+
.mockResolvedValueOnce(straddling)
1141+
.mockResolvedValue(CURRENT_ATTRIBUTION)
11411142
mockCheckAttributedUsageLimits.mockImplementation(
11421143
async (attribution: typeof CURRENT_ATTRIBUTION) => {
11431144
if (attribution.billingPeriod.end !== straddling.billingPeriod.end) {
@@ -1151,6 +1152,21 @@ describe('POST /api/billing/update-cost — mid-run usage gate', () => {
11511152
const body = await (await POST(attributedCallback())).json()
11521153

11531154
expect(body.usageExceeded).toBe(false)
1155+
expect(mockRefreshAttributionPeriod).toHaveBeenCalledTimes(2)
1156+
})
1157+
1158+
it('answers not exceeded when the standing read outlasts the callback budget', async () => {
1159+
mockCheckAttributedUsageLimits.mockImplementation(async () => {
1160+
await sleep(1500)
1161+
return { isExceeded: true, scope: 'payer' }
1162+
})
1163+
const startedAt = Date.now()
1164+
1165+
const res = await POST(attributedCallback())
1166+
1167+
expect(res.status).toBe(200)
1168+
await expect(res.json()).resolves.toMatchObject({ success: true, usageExceeded: false })
1169+
expect(Date.now() - startedAt).toBeLessThan(1400)
11541170
})
11551171

11561172
it('reloads a cached current period once it has ended', async () => {

‎apps/sim/app/api/billing/update-cost/route.ts‎

Lines changed: 22 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ import {
2222
resolveLegacyV0BillingAttribution,
2323
toBillingContext,
2424
} from '@/lib/billing/core/billing-attribution'
25-
import { readMidRunUsageVerdict } from '@/lib/billing/core/mid-run-usage'
25+
import { type MidRunUsageVerdict, readMidRunUsageVerdict } from '@/lib/billing/core/mid-run-usage'
2626
import {
2727
type CumulativeUsageContextField,
2828
CumulativeUsageContextMismatchError,
@@ -35,6 +35,7 @@ import {
3535
} from '@/lib/billing/threshold-billing'
3636
import { resolveUsageUpgradePayload } from '@/lib/billing/usage-upgrade'
3737
import { isBillingEnabled, isHosted } from '@/lib/core/config/env-flags'
38+
import { withinDeadline } from '@/lib/core/utils/deadline'
3839
import { generateRequestId } from '@/lib/core/utils/request'
3940
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
4041
import { BILLING_CALLBACK_OUTCOME } from '@/lib/mothership/generated/billing-protocol-v1'
@@ -45,6 +46,14 @@ import { checkInternalApiKey } from '@/lib/mothership/request/http'
4546
import { withIncomingGoSpan } from '@/lib/mothership/request/otel'
4647

4748
const logger = createLogger('BillingUpdateCostAPI')
49+
/**
50+
* How long a cost callback waits on the payer's standing. The worker gives up on the whole
51+
* callback after 5 s, and a cold gate read can wait on the ledger far longer; past this the
52+
* callback answers not-exceeded. The abandoned read keeps running and caches its admission, and
53+
* the next step or re-check reads a refusal again.
54+
*/
55+
const USAGE_STANDING_TIMEOUT_MS = 1000
56+
4857
const RETRYABLE_SETTLEMENT_RESPONSE = {
4958
code: 'BILLING_SETTLEMENT_RETRYABLE',
5059
error: 'Billing settlement temporarily unavailable',
@@ -70,14 +79,24 @@ function invalidBillingProtocolResponse(requestId: string, span: Span): NextResp
7079
* Served from the execution usage gate: an admission is cached per payer and actor for the gate
7180
* TTL and a refusal is always re-read, so steady-state steps cost no ledger read. The charge is
7281
* already recorded when this runs; a gate that cannot answer reports not-exceeded and leaves the
73-
* refusal to the next step or re-check rather than ending a paying run on a database blip.
82+
* refusal to the next step or re-check rather than ending a paying run on a database blip,
83+
* and so does a read that outlasts {@link USAGE_STANDING_TIMEOUT_MS}.
7484
*/
7585
async function readUsageStanding(
7686
userId: string,
7787
billingAttribution: BillingAttributionSnapshot | undefined
7888
): Promise<BillingUsageVerdict> {
7989
if (!isHosted || !billingAttribution) return { usageExceeded: false }
80-
const verdict = await readMidRunUsageVerdict(billingAttribution)
90+
let verdict: MidRunUsageVerdict
91+
try {
92+
verdict = await withinDeadline(
93+
() => readMidRunUsageVerdict(billingAttribution),
94+
Date.now() + USAGE_STANDING_TIMEOUT_MS
95+
)
96+
} catch {
97+
logger.warn('Usage standing read outlasted the callback budget; answering not exceeded')
98+
return { usageExceeded: false }
99+
}
81100
// Only a spent limit pauses the run. A blocked account is refused at the run's next
82101
// continuation or re-check, with blocked-account copy rather than the upgrade card.
83102
if (verdict.status !== 'exceeded') return { usageExceeded: false }

‎apps/sim/app/api/copilot/api-keys/validate/route.test.ts‎

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -141,7 +141,7 @@ vi.mock('@/lib/workspaces/permissions/utils', () => permissionsMock)
141141
vi.mock('@/lib/workspaces/utils', () => workspacesUtilsMock)
142142

143143
import { validateCopilotApiKeyBodySchema } from '@/lib/api/contracts/copilot'
144-
import { resetMidRunPeriodCache } from '@/lib/billing/core/mid-run-usage'
144+
import { resetMidRunUsageCaches } from '@/lib/billing/core/mid-run-usage'
145145
import { resetUsageGateCache } from '@/lib/billing/core/usage-gate-cache'
146146
import { POST } from '@/app/api/copilot/api-keys/validate/route'
147147

@@ -524,7 +524,7 @@ describe('validation lifecycle purposes', () => {
524524
periodEnd: new Date(ATTRIBUTION.billingPeriod.end),
525525
})
526526
resetUsageGateCache()
527-
resetMidRunPeriodCache()
527+
resetMidRunUsageCaches()
528528
})
529529

530530
it('defaults older callers to full admission and rejects unknown purposes', () => {
@@ -796,6 +796,28 @@ describe('validation lifecycle purposes', () => {
796796
expect(response.status).toBe(402)
797797
})
798798

799+
it('answers repeated direct-v1 continuations from the cached admission and re-reads a refusal', async () => {
800+
for (let call = 0; call < 2; call++) queueTableRows(schemaMock.user, [{ id: 'user-1' }])
801+
for (let leg = 0; leg < 3; leg++) {
802+
expect((await POST(request(body, directHeaders))).status).toBe(200)
803+
}
804+
expect(mockCheckUsageStatus).toHaveBeenCalledTimes(1)
805+
806+
resetMidRunUsageCaches()
807+
mockCheckUsageStatus.mockResolvedValue({ isExceeded: true, currentUsage: 12, limit: 10 })
808+
expect((await POST(request(body, directHeaders))).status).toBe(402)
809+
expect((await POST(request(body, directHeaders))).status).toBe(402)
810+
expect(mockCheckUsageStatus).toHaveBeenCalledTimes(3)
811+
})
812+
813+
it('never reads the ledger for a direct-v1 continuation when billing is off', async () => {
814+
setEnvFlags({ isHosted: false, isBillingEnabled: false })
815+
mockCheckUsageStatus.mockResolvedValue({ isExceeded: true, currentUsage: 12, limit: 10 })
816+
817+
expect((await POST(request(body, directHeaders))).status).toBe(200)
818+
expect(mockCheckUsageStatus).not.toHaveBeenCalled()
819+
})
820+
799821
it('judges a direct-v1 organization payer without a subscription as that organization', async () => {
800822
mockGetOrganizationSubscription.mockResolvedValue(null)
801823
mockCheckUsageStatus.mockImplementation(

‎apps/sim/lib/billing/core/mid-run-usage.ts‎

Lines changed: 49 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,13 @@ import {
1212
import { defaultBillingPeriod } from '@/lib/billing/core/billing-period'
1313
import { getHighestPriorityPersonalSubscription } from '@/lib/billing/core/plan'
1414
import { resolveSubscriptionUsagePeriod } from '@/lib/billing/core/reporting-period'
15-
import { checkExecutionUsageLimits } from '@/lib/billing/core/usage-gate-cache'
15+
import {
16+
checkExecutionUsageLimits,
17+
USAGE_GATE_SETTLE_TIMEOUT_MS,
18+
USAGE_GATE_TTL_MS,
19+
} from '@/lib/billing/core/usage-gate-cache'
20+
import { coalesceLocally } from '@/lib/concurrency/singleflight'
21+
import { isBillingEnabled, isHosted } from '@/lib/core/config/env-flags'
1622

1723
const logger = createLogger('MidRunUsage')
1824

@@ -121,6 +127,17 @@ export async function readMidRunUsageVerdict(
121127
return { status: 'unknown' }
122128
}
123129

130+
/**
131+
* Admitted direct-v1 verdicts, served for the execution gate's TTL like
132+
* {@link checkExecutionUsageLimits} serves attributed ones: the worker re-validates a run on
133+
* every resume leg, and each uncached read sums the payer's ledger for the period. Only a
134+
* `within` verdict is stored, so a refusal or an unreadable ledger is always read again.
135+
*/
136+
const accountVerdictCache = new LRUCache<string, MidRunUsageVerdict>({
137+
max: 10_000,
138+
ttl: USAGE_GATE_TTL_MS,
139+
})
140+
124141
/**
125142
* The same verdict for a direct-v1 run billed to an account decision rather than an attributed
126143
* payer. The payer is the one saved in the decision at admission, never re-selected from the
@@ -129,6 +146,7 @@ export async function readMidRunUsageVerdict(
129146
export async function readMidRunAccountUsageVerdict(
130147
decision: AccountBillingDecision
131148
): Promise<MidRunUsageVerdict> {
149+
if (!isHosted || !isBillingEnabled) return { status: 'within' }
132150
try {
133151
const payer = decision.billingEntity
134152
const subscription =
@@ -139,6 +157,20 @@ export async function readMidRunAccountUsageVerdict(
139157
...defaultBillingPeriod(),
140158
source: 'default' as const,
141159
}
160+
const key = [
161+
payer.type,
162+
payer.id,
163+
billingPeriod.start.toISOString(),
164+
billingPeriod.end.toISOString(),
165+
billingPeriod.source,
166+
decision.userId,
167+
subscription?.id ?? '',
168+
subscription?.plan ?? '',
169+
subscription?.status ?? '',
170+
subscription?.seats ?? '',
171+
].join(':')
172+
const cached = accountVerdictCache.get(key)
173+
if (cached) return cached
142174
// An organization payer without a subscription stays organization-scoped on the free plan,
143175
// as `toUsageLimitSubscription` does for attributed runs, never the actor's personal ledger.
144176
const usageSubscription =
@@ -153,12 +185,20 @@ export async function readMidRunAccountUsageVerdict(
153185
periodEnd: billingPeriod.end,
154186
}
155187
: null)
156-
const usage = await checkUsageStatus(decision.userId, usageSubscription, {
157-
billingEntity: payer,
158-
billingPeriod,
159-
})
188+
const usage = await coalesceLocally(
189+
`mid-run-account-usage:${key}`,
190+
() =>
191+
checkUsageStatus(decision.userId, usageSubscription, {
192+
billingEntity: payer,
193+
billingPeriod,
194+
}),
195+
USAGE_GATE_SETTLE_TIMEOUT_MS
196+
)
160197
if (usage.unavailable) return { status: 'unknown' }
161-
return usage.isExceeded ? { status: 'exceeded', scope: 'payer' } : { status: 'within' }
198+
if (usage.isExceeded) return { status: 'exceeded', scope: 'payer' }
199+
const within: MidRunUsageVerdict = { status: 'within' }
200+
accountVerdictCache.set(key, within)
201+
return within
162202
} catch (error) {
163203
logger.warn('Mid-run account usage read failed; continuing the run', {
164204
error: getErrorMessage(error),
@@ -167,7 +207,8 @@ export async function readMidRunAccountUsageVerdict(
167207
}
168208
}
169209

170-
/** Drops every cached current period. Test seam; never called in production code. */
171-
export function resetMidRunPeriodCache(): void {
210+
/** Drops every cached current period and account verdict. Test seam; never called in production code. */
211+
export function resetMidRunUsageCaches(): void {
172212
currentPeriodCache.clear()
213+
accountVerdictCache.clear()
173214
}

‎apps/sim/lib/mothership/request/lifecycle/admission.test.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import {
66
} from '@sim/testing/mocks/billing-usage-gate-cache.mock'
77
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
88
import { createAttributedBillingRequestEnvelope } from '@/lib/billing/core/billing-attribution'
9-
import { resetMidRunPeriodCache } from '@/lib/billing/core/mid-run-usage'
9+
import { resetMidRunUsageCaches } from '@/lib/billing/core/mid-run-usage'
1010
import { OrchestrationError } from '@/lib/core/orchestration/types'
1111
import { BillingLimitError } from '@/lib/mothership/request/go/stream'
1212
import { authorizeLifecycleContinuation, restoreBillingAdmission } from './admission'
@@ -49,7 +49,7 @@ beforeEach(() => {
4949
periodStart: new Date(attribution.billingPeriod.start),
5050
periodEnd: new Date(attribution.billingPeriod.end),
5151
})
52-
resetMidRunPeriodCache()
52+
resetMidRunUsageCaches()
5353
})
5454
afterEach(resetEnvFlagsMock)
5555

@@ -157,7 +157,7 @@ describe('continuation admission', () => {
157157
authorizeLifecycleContinuation({ ...context, billingAttribution: ended })
158158
).rejects.toBeInstanceOf(BillingLimitError)
159159

160-
resetMidRunPeriodCache()
160+
resetMidRunUsageCaches()
161161
mockGetOrganizationSubscription.mockRejectedValue(new Error('subscription read failed'))
162162
await expect(
163163
authorizeLifecycleContinuation({ ...context, billingAttribution: ended })

0 commit comments

Comments
 (0)