Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 8 additions & 9 deletions apps/sim/lib/billing/core/limit-notifications.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,6 @@ async function claimThreshold(
return writeLimitNotifications(scope, id, setExpr, onlyIfLower)
}

const DAY_MS = 24 * 60 * 60 * 1000

/** One account's credits threshold, keyed on its billing period and limit. */
export interface CreditsThresholdClaim {
scope: 'user' | 'organization'
Expand All @@ -92,20 +90,21 @@ export interface CreditsThresholdClaim {

/**
* The stored claims for a credits threshold, and the condition under which it is still unclaimed:
* `credits` holds the highest threshold emailed while `creditsPeriod` (the period's start day) and
* `creditsLimit` (the limit in cents) still match, so a new period or a changed limit — in either
* direction — re-arms both thresholds with no reset write, while within one a claim of 100 also
* retires 80, and never the reverse.
* `credits` holds the highest threshold emailed while `creditsPeriod` (the period's exact start,
* in epoch seconds) and `creditsLimit` (the limit in cents) still match, so a new period — even
* one starting the same day as the last — or a changed limit, in either direction, re-arms both
* thresholds with no reset write, while within one a claim of 100 also retires 80, and never the
* reverse.
*/
function creditsThresholdSql(claim: CreditsThresholdClaim) {
const periodDay = Math.floor(claim.periodStart.getTime() / DAY_MS)
const periodStartSeconds = Math.floor(claim.periodStart.getTime() / 1000)
const limitCents = Math.round(claim.limit * 100)
const column =
claim.scope === 'user' ? userStats.limitNotifications : organization.limitNotifications
return {
next: sql`coalesce(${column}, '{}'::jsonb) || jsonb_build_object('credits', ${claim.threshold}::int, 'creditsPeriod', ${periodDay}::bigint, 'creditsLimit', ${limitCents}::bigint)`,
next: sql`coalesce(${column}, '{}'::jsonb) || jsonb_build_object('credits', ${claim.threshold}::int, 'creditsPeriod', ${periodStartSeconds}::bigint, 'creditsLimit', ${limitCents}::bigint)`,
unclaimed: sql<boolean>`not (
(${column} ->> 'creditsPeriod')::bigint is not distinct from ${periodDay}::bigint
(${column} ->> 'creditsPeriod')::bigint is not distinct from ${periodStartSeconds}::bigint
and (${column} ->> 'creditsLimit')::bigint is not distinct from ${limitCents}::bigint
and coalesce((${column} ->> 'credits')::int, 0) >= ${claim.threshold}::int
)`,
Expand Down
39 changes: 38 additions & 1 deletion apps/sim/lib/billing/core/reporting-usage-cache.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,10 @@ const redisUrl = readTestRedisUrl()
vi.mock('@sim/db', () => ({ db: { transaction }, dbReplica: {} }))
vi.mock('@/lib/core/config/redis', () => redisConfigMock)

import { readSoftGateUsageCost } from '@/lib/billing/core/reporting-usage-cache'
import {
REPORTING_USAGE_CACHE_TTL_MS,
readSoftGateUsageCost,
} from '@/lib/billing/core/reporting-usage-cache'
import type { BillingEntity, UsageQueryPeriod } from '@/lib/billing/core/usage-log'

const schemaName = `reporting_usage_${generateId().replaceAll('-', '')}`
Expand Down Expand Up @@ -71,6 +74,7 @@ describe.runIf(Boolean(redisUrl))('shared reporting usage read', () => {
})

afterEach(async () => {
vi.useRealTimers()
await redis.del(sharedKey(payer))
})

Expand Down Expand Up @@ -101,6 +105,39 @@ describe.runIf(Boolean(redisUrl))('shared reporting usage read', () => {
await expect(readSoftGateUsageCost(payer, REPORTING)).resolves.toBe(5.75)
})

it('shares a sum that ran longer than the TTL for a short floor instead of dropping it', async () => {
const now = Date.now()
vi.useFakeTimers({ toFake: ['Date'] })
vi.setSystemTime(now - REPORTING_USAGE_CACHE_TTL_MS - 1_000)
transaction.mockImplementationOnce(async (callback) => {
const result = await database.transaction(callback)
vi.setSystemTime(now)
return result
})

await expect(readSoftGateUsageCost(payer, REPORTING)).resolves.toBe(5.75)
expect(await redis.get(sharedKey(payer))).toBe('5.75')
const ttl = await redis.pttl(sharedKey(payer))
expect(ttl).toBeGreaterThan(4_000)
expect(ttl).toBeLessThanOrEqual(5_000)
})

it('anchors a stored sum expiry to when the sum began, not when it was written', async () => {
const now = Date.now()
vi.useFakeTimers({ toFake: ['Date'] })
vi.setSystemTime(now - 20_000)
transaction.mockImplementationOnce(async (callback) => {
const result = await database.transaction(callback)
vi.setSystemTime(now)
return result
})

await expect(readSoftGateUsageCost(payer, REPORTING)).resolves.toBe(5.75)
const ttl = await redis.pttl(sharedKey(payer))
expect(ttl).toBeGreaterThan(REPORTING_USAGE_CACHE_TTL_MS - 20_000 - 1_000)
expect(ttl).toBeLessThanOrEqual(REPORTING_USAGE_CACHE_TTL_MS - 20_000 + 5_000)
})

it('never lets a slower, older sum replace one stored while it ran', async () => {
transaction.mockImplementationOnce(async (callback) => {
const result = await database.transaction(callback)
Expand Down
40 changes: 30 additions & 10 deletions apps/sim/lib/billing/core/reporting-usage-cache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,15 +28,24 @@ const logger = createLogger('ReportingUsageCache')
* small against a year-long allowance while turning a per-event scan into one per window.
*
* A sum is held both in Redis, shared by every process, and in each process that reads it. It
* reflects the ledger as of the moment its sum began, so a served sum can omit usage written over
* the sum's own duration, plus up to this TTL and its jitter in Redis, plus up to this TTL again
* in the reading process.
* reflects the ledger as of the moment its sum began, and its Redis expiry is anchored to that
* moment, so it is served from Redis for at most max(this TTL + jitter, the sum's duration +
* {@link MIN_SHARED_TTL_MS}) after it began (plus any reconnect delay for a resent write), plus
* up to this TTL again in the reading process.
*/
export const REPORTING_USAGE_CACHE_TTL_MS = 30_000

/** Redis expiry is jittered by up to this much, so payers summed together do not expire together. */
const SHARED_TTL_JITTER_MS = 5_000

/**
* The shortest time a sum is kept in Redis. A sum that ran longer than the TTL would otherwise
* expire on arrival, and under the database pressure that makes sums slow, every process would
* then run the same slow sum again. Five seconds shares it with the processes waiting on it while
* adding little staleness next to the time the sum itself took.
*/
const MIN_SHARED_TTL_MS = 5_000

/**
* How long a read waits on a connected Redis before summing the ledger instead, so a socket that
* has silently stopped answering costs one sum rather than the shared client's long timeouts.
Expand Down Expand Up @@ -95,17 +104,27 @@ function warnSharedWriteFailed(error: unknown): void {
}

/**
* Fire-and-forget: a read never waits on, or fails because of, the shared write. The write only
* lands when no sum is stored (`NX`), so a slow, older sum can never replace a fresher one or
* extend its expiry.
* Fire-and-forget: a read never waits on, or fails because of, the shared write. The expiry is
* anchored to when the sum began: the remaining lifetime is computed here and written with `PX`,
* which every Redis version accepts, floored at {@link MIN_SHARED_TTL_MS}. A write the client
* resends after a reconnect re-applies that same relative lifetime from the resend, so it can
* extend the expiry by at most the reconnect delay. The write only lands when no sum is stored
* (`NX`), so an older sum can never replace a fresher one or extend its expiry.
*/
function writeSharedReportingUsageCost(key: string, cost: number): void {
function writeSharedReportingUsageCost(key: string, cost: number, sumStartedAt: number): void {
try {
const redis = readyRedisClient()
if (!redis) return
const ttlMs = REPORTING_USAGE_CACHE_TTL_MS + randomInt(0, SHARED_TTL_JITTER_MS)
const remainingMs =
sumStartedAt + REPORTING_USAGE_CACHE_TTL_MS + randomInt(0, SHARED_TTL_JITTER_MS) - Date.now()
redis
.set(sharedReportingUsageKey(key), String(cost), 'PX', ttlMs, 'NX')
.set(
sharedReportingUsageKey(key),
String(cost),
'PX',
Math.max(remainingMs, MIN_SHARED_TTL_MS),
'NX'
)
.catch(warnSharedWriteFailed)
} catch (error) {
warnSharedWriteFailed(error)
Expand All @@ -124,8 +143,9 @@ async function sumReportingUsageCost(
): Promise<number> {
const shared = await readSharedReportingUsageCost(key)
if (shared !== undefined) return shared
const sumStartedAt = Date.now()
const cost = await getBillingPeriodUsageCost(entity, period)
writeSharedReportingUsageCost(key, cost)
writeSharedReportingUsageCost(key, cost, sumStartedAt)
return cost
}

Expand Down
11 changes: 10 additions & 1 deletion apps/sim/lib/billing/core/usage-threshold-email.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ async function claims(): Promise<Record<string, number>> {
function claimOf(threshold: 80 | 100, periodStart = SEPTEMBER, limitCents = 10_000) {
return {
credits: threshold,
creditsPeriod: Math.floor(periodStart.getTime() / 86_400_000),
creditsPeriod: periodStart.getTime() / 1000,
creditsLimit: limitCents,
}
}
Expand Down Expand Up @@ -159,6 +159,15 @@ describe('usage threshold email', () => {
expect(await claims()).toEqual(claimOf(80))
})

it('re-arms for a new period that starts the same day as the one it replaces', async () => {
const replacement = new Date('2026-09-01T12:00:00.000Z')
await notify(85)
await notify(85, { periodStart: replacement })

expect(delivered()).toEqual([warning(), warning()])
expect(await claims()).toEqual(claimOf(80, replacement))
})

it('warns again at a raised limit after the old one was reached', async () => {
await notify(100)
await notify(100, { limit: 125 })
Expand Down
5 changes: 3 additions & 2 deletions packages/db/schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1359,8 +1359,9 @@ export const userStats = pgTable('user_stats', {
* re-arms when usage drops back below the re-arm band. Keyed by limit
* category ('storage' | 'tables'); seats live on `organization`. `credits`
* instead holds the threshold emailed for the billing period and limit in
* `creditsPeriod` (start day) and `creditsLimit` (cents), so a new period or a
* changed limit re-arms it without a reset (see `claimCreditsThreshold`).
* `creditsPeriod` (start, epoch seconds) and `creditsLimit` (cents), so a new
* period or a changed limit re-arms it without a reset (see
* `claimCreditsThreshold`).
*
* Dedup granularity is per billing account per category — intentionally NOT
* per table, so a user hitting the row limit on several tables gets one
Expand Down
Loading