diff --git a/apps/sim/lib/billing/core/limit-notifications.ts b/apps/sim/lib/billing/core/limit-notifications.ts index 30c2238aee2..8a4311e2811 100644 --- a/apps/sim/lib/billing/core/limit-notifications.ts +++ b/apps/sim/lib/billing/core/limit-notifications.ts @@ -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' @@ -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`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 )`, diff --git a/apps/sim/lib/billing/core/reporting-usage-cache.integration.ts b/apps/sim/lib/billing/core/reporting-usage-cache.integration.ts index 3efd5100633..a8986e95741 100644 --- a/apps/sim/lib/billing/core/reporting-usage-cache.integration.ts +++ b/apps/sim/lib/billing/core/reporting-usage-cache.integration.ts @@ -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('-', '')}` @@ -71,6 +74,7 @@ describe.runIf(Boolean(redisUrl))('shared reporting usage read', () => { }) afterEach(async () => { + vi.useRealTimers() await redis.del(sharedKey(payer)) }) @@ -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) diff --git a/apps/sim/lib/billing/core/reporting-usage-cache.ts b/apps/sim/lib/billing/core/reporting-usage-cache.ts index 440b69ca2dc..56ac65d00d0 100644 --- a/apps/sim/lib/billing/core/reporting-usage-cache.ts +++ b/apps/sim/lib/billing/core/reporting-usage-cache.ts @@ -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. @@ -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) @@ -124,8 +143,9 @@ async function sumReportingUsageCost( ): Promise { 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 } diff --git a/apps/sim/lib/billing/core/usage-threshold-email.integration.ts b/apps/sim/lib/billing/core/usage-threshold-email.integration.ts index c78c90a4c3c..366c7e2bff1 100644 --- a/apps/sim/lib/billing/core/usage-threshold-email.integration.ts +++ b/apps/sim/lib/billing/core/usage-threshold-email.integration.ts @@ -92,7 +92,7 @@ async function claims(): Promise> { 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, } } @@ -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 }) diff --git a/packages/db/schema.ts b/packages/db/schema.ts index 5766d192b85..a8669b34722 100644 --- a/packages/db/schema.ts +++ b/packages/db/schema.ts @@ -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