mirror of
https://github.com/simstudioai/sim.git
synced 2026-09-24 15:45:35 +08:00
feat(billing): unify upgrade routing with reason context + storage/tables limit emails (#5171)
* feat(billing): unify upgrade routing with reason context + storage/tables limit emails * fix(billing): re-arm limit-notification dedup on usage drops (prior-usage + decrement) * fix(billing): isolate per-admin email failures in org limit notifications * fix(billing): re-arm limit dedup at zero usage and zero prior usage (full clear / wipe-rebuild) * fix(billing): make storage-decrement notification re-arm only (never send on a shrink) * fix(billing): resolve recipients before claiming so opt-outs don't burn the dedup threshold * fix(billing): fire table limit emails on upsert inserts via shared notifyTableRowUsage * chore(billing): only log a limit email as sent when a recipient actually received it * chore(billing): match to_jsonb int cast between claim and re-arm for consistency * fix(billing): notify table limits post-commit so a rolled-back insert never emails or burns the claim * feat(pi): swap Pi Coding Agent icon to the pi glyph and use a black bgColor * fix(billing): drop priorUsage re-arm to make dedup a single atomic claim (no duplicate-email race) * docs(billing): move limit-notification rationale to TSDoc, correct tables warn-once behavior * docs(db): note limit_notifications dedup is per-account, not per-table * perf(billing): cut redundant subscription fetches and edge-gate notify to slash DB load * docs(billing): drop self-explanatory inline comments from the notification path
This commit is contained in:
@@ -5318,6 +5318,18 @@ export function SmtpIcon(props: SVGProps<SVGSVGElement>) {
|
||||
)
|
||||
}
|
||||
|
||||
export function PiIcon(props: SVGProps<SVGSVGElement>) {
|
||||
return (
|
||||
<svg {...props} xmlns='http://www.w3.org/2000/svg' viewBox='0 0 800 800' fill='currentColor'>
|
||||
<path
|
||||
fillRule='evenodd'
|
||||
d='M165.29 165.29 H517.36 V400 H400 V517.36 H282.65 V634.72 H165.29 Z M282.65 282.65 V400 H400 V282.65 Z'
|
||||
/>
|
||||
<path d='M517.36 400 H634.72 V634.72 H517.36 Z' />
|
||||
</svg>
|
||||
)
|
||||
}
|
||||
|
||||
export function SshIcon(props: SVGProps<SVGSVGElement>) {
|
||||
return (
|
||||
<svg
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
|
||||
import { type DragEvent, useCallback, useEffect, useMemo, useRef, useState } from 'react'
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { toError } from '@sim/utils/errors'
|
||||
import { getErrorMessage, toError } from '@sim/utils/errors'
|
||||
import { useParams, useRouter } from 'next/navigation'
|
||||
import { useQueryStates } from 'nuqs'
|
||||
import { usePostHog } from 'posthog-js/react'
|
||||
@@ -25,6 +25,7 @@ import {
|
||||
} from '@/components/emcn'
|
||||
import { Download, Send } from '@/components/emcn/icons'
|
||||
import { getDocumentIcon } from '@/components/icons/document-icons'
|
||||
import { useLimitUpgradeToast } from '@/lib/billing/client'
|
||||
import { captureEvent } from '@/lib/posthog/client'
|
||||
import { triggerFileDownload } from '@/lib/uploads/client/download'
|
||||
import type { WorkspaceFileRecord } from '@/lib/uploads/contexts/workspace'
|
||||
@@ -197,6 +198,7 @@ export function Files() {
|
||||
const { data: folders = EMPTY_WORKSPACE_FILE_FOLDERS } = useWorkspaceFileFolders(workspaceId)
|
||||
const { data: members } = useWorkspaceMembersQuery(workspaceId)
|
||||
const uploadFile = useUploadWorkspaceFile()
|
||||
const notifyLimit = useLimitUpgradeToast()
|
||||
const deleteFile = useDeleteWorkspaceFile()
|
||||
const renameFile = useRenameWorkspaceFile()
|
||||
const createFolder = useCreateWorkspaceFileFolder()
|
||||
@@ -699,6 +701,12 @@ export function Files() {
|
||||
})
|
||||
} catch (err) {
|
||||
logger.error('Error uploading file:', err)
|
||||
const message = getErrorMessage(err)
|
||||
if (/storage limit/i.test(message)) {
|
||||
notifyLimit('storage', message)
|
||||
} else {
|
||||
toast.error(`Failed to upload "${allowedFiles[i].name}"`)
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
@@ -708,7 +716,7 @@ export function Files() {
|
||||
setUploadProgress({ completed: 0, total: 0, currentPercent: 0 })
|
||||
}
|
||||
},
|
||||
[workspaceId, canEdit, currentFolderId]
|
||||
[workspaceId, canEdit, currentFolderId, notifyLimit]
|
||||
)
|
||||
|
||||
const rowDragDropConfig = useMemo<RowDragDropConfig>(
|
||||
|
||||
@@ -7,6 +7,7 @@ import { Chip } from '@/components/emcn'
|
||||
import { Credit } from '@/components/emcn/icons'
|
||||
import { ON_DEMAND_UNLIMITED } from '@/lib/billing/constants'
|
||||
import { formatCredits } from '@/lib/billing/credits/conversion'
|
||||
import { buildUpgradeHref } from '@/lib/billing/upgrade-reasons'
|
||||
import { isBillingEnabled } from '@/app/workspace/[workspaceId]/settings/navigation'
|
||||
import { useMyMemberCredits } from '@/hooks/queries/organization'
|
||||
import { usePlanView } from '@/hooks/queries/plan-view'
|
||||
@@ -33,7 +34,7 @@ function CreditsChipInner() {
|
||||
const { workspaceId } = useParams<{ workspaceId: string }>()
|
||||
const { data: memberCredits, isLoading: memberLoading } = useMyMemberCredits(workspaceId)
|
||||
|
||||
const upgradeHref = `/workspace/${workspaceId}/upgrade`
|
||||
const upgradeHref = buildUpgradeHref(workspaceId, 'credits')
|
||||
|
||||
/**
|
||||
* Warm the route bundle and the exact queries the Upgrade page gates on, so
|
||||
|
||||
@@ -33,6 +33,7 @@ import {
|
||||
hasPaidSubscriptionStatus,
|
||||
hasUsableSubscriptionAccess,
|
||||
} from '@/lib/billing/subscriptions/utils'
|
||||
import { buildUpgradeHref } from '@/lib/billing/upgrade-reasons'
|
||||
import { cn } from '@/lib/core/utils/cn'
|
||||
import { getBaseUrl } from '@/lib/core/utils/urls'
|
||||
import { UsageLimitField } from '@/app/workspace/[workspaceId]/settings/components/billing/components/usage-limit-field/usage-limit-field'
|
||||
@@ -125,7 +126,7 @@ export function Billing() {
|
||||
const betterAuthSubscription = useSubscription()
|
||||
const openBillingPortal = useOpenBillingPortal()
|
||||
|
||||
const upgradeHref = `/workspace/${workspaceId}/upgrade`
|
||||
const upgradeHref = buildUpgradeHref(workspaceId)
|
||||
|
||||
/**
|
||||
* Warm the Upgrade route bundle and the exact queries that page gates on, so
|
||||
|
||||
@@ -24,6 +24,7 @@ import {
|
||||
workspaceRoleLockReason,
|
||||
} from '@/components/permissions'
|
||||
import type { WorkspacePermission } from '@/lib/api/contracts/workspaces'
|
||||
import { buildUpgradeHref } from '@/lib/billing/upgrade-reasons'
|
||||
import {
|
||||
MemberRow,
|
||||
MemberSection,
|
||||
@@ -105,7 +106,7 @@ export function Teammates() {
|
||||
const inviteDisabledReason = activeWorkspace?.inviteDisabledReason ?? null
|
||||
const isInvitationsDisabled = isInvitationsDisabledByConfig || inviteDisabledReason !== null
|
||||
|
||||
const upgradeHref = `/workspace/${workspaceId}/upgrade`
|
||||
const upgradeHref = buildUpgradeHref(workspaceId, 'seats')
|
||||
|
||||
/**
|
||||
* Warm the Upgrade route bundle and the queries it gates on, so a gated
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import { Suspense } from 'react'
|
||||
import type { Metadata } from 'next'
|
||||
import { Upgrade } from '@/app/workspace/[workspaceId]/upgrade/upgrade'
|
||||
|
||||
@@ -9,5 +10,9 @@ export default async function UpgradePage({
|
||||
params: Promise<{ workspaceId: string }>
|
||||
}) {
|
||||
const { workspaceId } = await params
|
||||
return <Upgrade workspaceId={workspaceId} />
|
||||
return (
|
||||
<Suspense fallback={<div className='h-full bg-[var(--bg)]' />}>
|
||||
<Upgrade workspaceId={workspaceId} />
|
||||
</Suspense>
|
||||
)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
import { parseAsStringLiteral } from 'nuqs/server'
|
||||
import { UPGRADE_REASONS } from '@/lib/billing/upgrade-reasons'
|
||||
|
||||
/**
|
||||
* Single source of truth for the upgrade page's `reason` query param.
|
||||
*
|
||||
* Nullable (no `.withDefault`): a clean URL means no reason and the page keeps
|
||||
* its generic header. Shared by the client (`useQueryState`) and any server
|
||||
* read via `createSearchParamsCache`.
|
||||
*/
|
||||
export const upgradeReasonParam = {
|
||||
key: 'reason',
|
||||
parser: parseAsStringLiteral(UPGRADE_REASONS),
|
||||
} as const
|
||||
|
||||
/** Clean URLs, no back-stack churn — the reason is a passive header hint. */
|
||||
export const upgradeUrlKeys = {
|
||||
history: 'replace',
|
||||
clearOnDefault: true,
|
||||
} as const
|
||||
@@ -3,6 +3,7 @@
|
||||
import { useCallback, useEffect, useState } from 'react'
|
||||
import { getErrorMessage } from '@sim/utils/errors'
|
||||
import { useRouter } from 'next/navigation'
|
||||
import { useQueryState } from 'nuqs'
|
||||
import { ArrowLeft, Chip, toast } from '@/components/emcn'
|
||||
import {
|
||||
getUpgradeCardCta,
|
||||
@@ -11,6 +12,7 @@ import {
|
||||
type UpgradeCardId,
|
||||
} from '@/lib/billing/client'
|
||||
import { ANNUAL_DISCOUNT_RATE } from '@/lib/billing/constants'
|
||||
import { DEFAULT_UPGRADE_HEADER, UPGRADE_REASON_COPY } from '@/lib/billing/upgrade-reasons'
|
||||
import { isBillingEnabled } from '@/app/workspace/[workspaceId]/settings/navigation'
|
||||
import {
|
||||
BillingPeriodToggle,
|
||||
@@ -26,6 +28,10 @@ import {
|
||||
PRO_PLAN_CREDITS,
|
||||
PRO_PLAN_FEATURES,
|
||||
} from '@/app/workspace/[workspaceId]/upgrade/plan-configs'
|
||||
import {
|
||||
upgradeReasonParam,
|
||||
upgradeUrlKeys,
|
||||
} from '@/app/workspace/[workspaceId]/upgrade/search-params'
|
||||
import { useFullscreenOriginStore } from '@/stores/fullscreen-origin'
|
||||
|
||||
const TYPEFORM_ENTERPRISE_URL = 'https://form.typeform.com/to/jqCO12pF' as const
|
||||
@@ -47,8 +53,14 @@ export function Upgrade({ workspaceId }: UpgradeProps) {
|
||||
const state = useUpgradeState()
|
||||
const router = useRouter()
|
||||
const origin = useFullscreenOriginStore((s) => s.origin)
|
||||
const [reason] = useQueryState(upgradeReasonParam.key, {
|
||||
...upgradeReasonParam.parser,
|
||||
...upgradeUrlKeys,
|
||||
})
|
||||
const [showAllFeatures, setShowAllFeatures] = useState(false)
|
||||
|
||||
const header = reason ? UPGRADE_REASON_COPY[reason].header : DEFAULT_UPGRADE_HEADER
|
||||
|
||||
const handleBack = useCallback(() => {
|
||||
router.replace(origin ?? `/workspace/${workspaceId}/home`)
|
||||
}, [origin, router, workspaceId])
|
||||
@@ -152,7 +164,7 @@ export function Upgrade({ workspaceId }: UpgradeProps) {
|
||||
<div className='mx-auto flex w-full max-w-[960px] flex-col gap-7 pt-6 pb-3'>
|
||||
<div className='flex flex-col items-center gap-4'>
|
||||
<h1 className='text-balance text-center font-season text-[30px] text-[var(--text-primary)]'>
|
||||
Plans that scale with you
|
||||
{header}
|
||||
</h1>
|
||||
{state.showUpgradePlans && (
|
||||
<BillingPeriodToggle isAnnual={state.isAnnual} onChange={state.setIsAnnual} />
|
||||
|
||||
+2
-1
@@ -4,6 +4,7 @@ import { useQueryClient } from '@tanstack/react-query'
|
||||
import { ArrowRight } from 'lucide-react'
|
||||
import { useParams, useRouter } from 'next/navigation'
|
||||
import { ChipLink } from '@/components/emcn'
|
||||
import { buildUpgradeHref } from '@/lib/billing/upgrade-reasons'
|
||||
import { prefetchUpgradeBillingData } from '@/hooks/queries/subscription'
|
||||
import { prefetchWorkspaceSettings } from '@/hooks/queries/workspace'
|
||||
|
||||
@@ -15,7 +16,7 @@ export function DeployUpgradeGate({ feature }: DeployUpgradeGateProps) {
|
||||
const router = useRouter()
|
||||
const queryClient = useQueryClient()
|
||||
const { workspaceId } = useParams<{ workspaceId: string }>()
|
||||
const upgradeHref = `/workspace/${workspaceId}/upgrade`
|
||||
const upgradeHref = buildUpgradeHref(workspaceId)
|
||||
|
||||
// Warm the upgrade route + the queries it gates on so the click lands on
|
||||
// cached data. ChipLink isn't memoized, so no useCallback is needed.
|
||||
|
||||
@@ -53,7 +53,7 @@ export const PiBlock: BlockConfig<PiResponse> = {
|
||||
`,
|
||||
category: 'blocks',
|
||||
integrationType: IntegrationType.AI,
|
||||
bgColor: '#6E56CF',
|
||||
bgColor: '#000000',
|
||||
icon: PiIcon,
|
||||
subBlocks: [
|
||||
{
|
||||
|
||||
@@ -3,6 +3,7 @@ export { CreditPurchaseEmail } from './credit-purchase-email'
|
||||
export { CreditsExhaustedEmail } from './credits-exhausted-email'
|
||||
export { EnterpriseSubscriptionEmail } from './enterprise-subscription-email'
|
||||
export { FreeTierUpgradeEmail } from './free-tier-upgrade-email'
|
||||
export { LimitThresholdEmail } from './limit-threshold-email'
|
||||
export { PaymentFailedEmail } from './payment-failed-email'
|
||||
export { PlanWelcomeEmail } from './plan-welcome-email'
|
||||
export { UsageThresholdEmail } from './usage-threshold-email'
|
||||
|
||||
@@ -0,0 +1,76 @@
|
||||
import { Link, Section, Text } from '@react-email/components'
|
||||
import { baseStyles } from '@/components/emails/_styles'
|
||||
import { EmailLayout } from '@/components/emails/components'
|
||||
import { UPGRADE_REASON_COPY, type UpgradeReason } from '@/lib/billing/upgrade-reasons'
|
||||
import { getBrandConfig } from '@/ee/whitelabeling'
|
||||
|
||||
interface LimitThresholdEmailProps {
|
||||
/** `warning` = approaching the limit (~80%); `reached` = at/over the limit. */
|
||||
kind: 'warning' | 'reached'
|
||||
/** Limit category, drives the shared copy. */
|
||||
reason: UpgradeReason
|
||||
userName?: string
|
||||
/** Pre-formatted current usage, e.g. "4.2 GB", "48,000 rows", "9 seats". */
|
||||
usageLabel: string
|
||||
/** Pre-formatted limit, e.g. "5 GB", "50,000 rows", "10 seats". */
|
||||
limitLabel: string
|
||||
percentUsed: number
|
||||
upgradeLink: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Single template for the per-category usage-limit emails (storage, tables,
|
||||
* seats). Copy comes from {@link UPGRADE_REASON_COPY} so the email language
|
||||
* matches the upgrade-page header the user lands on.
|
||||
*/
|
||||
export function LimitThresholdEmail({
|
||||
kind,
|
||||
reason,
|
||||
userName,
|
||||
usageLabel,
|
||||
limitLabel,
|
||||
percentUsed,
|
||||
upgradeLink,
|
||||
}: LimitThresholdEmailProps) {
|
||||
const brand = getBrandConfig()
|
||||
const copy = UPGRADE_REASON_COPY[reason]
|
||||
const lead = kind === 'reached' ? copy.reachedLead : copy.warningLead
|
||||
const previewText = `${brand.name}: ${lead}`
|
||||
|
||||
return (
|
||||
<EmailLayout preview={previewText} showUnsubscribe={true}>
|
||||
<Text style={{ ...baseStyles.paragraph, marginTop: 0 }}>
|
||||
{userName ? `Hi ${userName},` : 'Hi,'}
|
||||
</Text>
|
||||
|
||||
<Text style={baseStyles.paragraph}>
|
||||
{lead} Upgrade your plan for more {copy.noun}.
|
||||
</Text>
|
||||
|
||||
<Section style={baseStyles.infoBox}>
|
||||
<Text style={baseStyles.infoBoxTitle}>Usage</Text>
|
||||
<Text style={baseStyles.infoBoxList}>
|
||||
{usageLabel} of {limitLabel} used ({percentUsed}%)
|
||||
</Text>
|
||||
</Section>
|
||||
|
||||
{/* Divider */}
|
||||
<div style={baseStyles.divider} />
|
||||
|
||||
<Link href={upgradeLink} style={{ textDecoration: 'none' }}>
|
||||
<Text style={baseStyles.button}>Upgrade</Text>
|
||||
</Link>
|
||||
|
||||
{/* Divider */}
|
||||
<div style={baseStyles.divider} />
|
||||
|
||||
<Text style={{ ...baseStyles.footerText, textAlign: 'left' }}>
|
||||
{kind === 'reached'
|
||||
? 'One-time notification at 100% usage.'
|
||||
: 'One-time notification at 80% usage.'}
|
||||
</Text>
|
||||
</EmailLayout>
|
||||
)
|
||||
}
|
||||
|
||||
export default LimitThresholdEmail
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
CreditsExhaustedEmail,
|
||||
EnterpriseSubscriptionEmail,
|
||||
FreeTierUpgradeEmail,
|
||||
LimitThresholdEmail,
|
||||
PaymentFailedEmail,
|
||||
PlanWelcomeEmail,
|
||||
UsageThresholdEmail,
|
||||
@@ -24,9 +25,10 @@ import {
|
||||
WorkspaceInvitationEmail,
|
||||
} from '@/components/emails/invitations'
|
||||
import { HelpConfirmationEmail } from '@/components/emails/support'
|
||||
import type { UpgradeReason } from '@/lib/billing/upgrade-reasons'
|
||||
import { getBaseUrl } from '@/lib/core/utils/urls'
|
||||
|
||||
export { getEmailSubject } from './subjects'
|
||||
export { getEmailSubject, getLimitEmailSubject } from './subjects'
|
||||
|
||||
interface WorkspaceInvitation {
|
||||
workspaceId: string
|
||||
@@ -153,6 +155,18 @@ export async function renderFreeTierUpgradeEmail(params: {
|
||||
)
|
||||
}
|
||||
|
||||
export async function renderLimitThresholdEmail(params: {
|
||||
kind: 'warning' | 'reached'
|
||||
reason: UpgradeReason
|
||||
userName?: string
|
||||
usageLabel: string
|
||||
limitLabel: string
|
||||
percentUsed: number
|
||||
upgradeLink: string
|
||||
}): Promise<string> {
|
||||
return await render(LimitThresholdEmail(params))
|
||||
}
|
||||
|
||||
export async function renderPlanWelcomeEmail(params: {
|
||||
planName: string
|
||||
userName?: string
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import { UPGRADE_REASON_COPY, type UpgradeReason } from '@/lib/billing/upgrade-reasons'
|
||||
import { getBrandConfig } from '@/ee/whitelabeling'
|
||||
|
||||
/** Email subject type for all supported email templates */
|
||||
@@ -79,3 +80,15 @@ export function getEmailSubject(type: EmailSubjectType): string {
|
||||
return brandName
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Subject line for a per-category usage-limit email. Reuses the shared
|
||||
* {@link UPGRADE_REASON_COPY} so the subject matches the email body and the
|
||||
* upgrade-page header the user lands on.
|
||||
*/
|
||||
export function getLimitEmailSubject(reason: UpgradeReason, kind: 'warning' | 'reached'): string {
|
||||
const brandName = getBrandConfig().name
|
||||
const copy = UPGRADE_REASON_COPY[reason]
|
||||
const subject = kind === 'reached' ? copy.reachedSubject : copy.warningSubject
|
||||
return `${subject} on ${brandName}`
|
||||
}
|
||||
|
||||
@@ -5300,22 +5300,12 @@ export function SmtpIcon(props: SVGProps<SVGSVGElement>) {
|
||||
|
||||
export function PiIcon(props: SVGProps<SVGSVGElement>) {
|
||||
return (
|
||||
<svg
|
||||
{...props}
|
||||
xmlns='http://www.w3.org/2000/svg'
|
||||
width='24'
|
||||
height='24'
|
||||
viewBox='0 0 24 24'
|
||||
fill='none'
|
||||
stroke='currentColor'
|
||||
strokeWidth='2'
|
||||
strokeLinecap='round'
|
||||
strokeLinejoin='round'
|
||||
>
|
||||
<rect x='2' y='3' width='20' height='18' rx='2' />
|
||||
<path d='M7 9h10' />
|
||||
<path d='M10 9v7' />
|
||||
<path d='M15 9v7' />
|
||||
<svg {...props} xmlns='http://www.w3.org/2000/svg' viewBox='0 0 800 800' fill='currentColor'>
|
||||
<path
|
||||
fillRule='evenodd'
|
||||
d='M165.29 165.29 H517.36 V400 H400 V517.36 H282.65 V634.72 H165.29 Z M282.65 282.65 V400 H400 V282.65 Z'
|
||||
/>
|
||||
<path d='M517.36 400 H634.72 V634.72 H517.36 Z' />
|
||||
</svg>
|
||||
)
|
||||
}
|
||||
|
||||
@@ -71,6 +71,7 @@ import {
|
||||
updateTableRowContract,
|
||||
updateWorkflowGroupContract,
|
||||
} from '@/lib/api/contracts/tables'
|
||||
import { buildUpgradeHref } from '@/lib/billing/upgrade-reasons'
|
||||
import type {
|
||||
CsvHeaderMapping,
|
||||
EnrichmentRunDetail,
|
||||
@@ -694,7 +695,7 @@ export function useCreateTableRow({ workspaceId, tableId }: RowMutationContext)
|
||||
})
|
||||
},
|
||||
onError: (error) =>
|
||||
notifyRowWriteError(error, () => router.push(`/workspace/${workspaceId}/upgrade`)),
|
||||
notifyRowWriteError(error, () => router.push(buildUpgradeHref(workspaceId, 'tables'))),
|
||||
onSettled: () => {
|
||||
// `reconcileCreatedRow` (onSuccess) is the source of truth for the rows
|
||||
// cache + its `totalCount`; only refresh the count surfaces here so a late
|
||||
@@ -874,7 +875,7 @@ export function useBatchCreateTableRows({ workspaceId, tableId }: RowMutationCon
|
||||
})
|
||||
},
|
||||
onError: (error) =>
|
||||
notifyRowWriteError(error, () => router.push(`/workspace/${workspaceId}/upgrade`)),
|
||||
notifyRowWriteError(error, () => router.push(buildUpgradeHref(workspaceId, 'tables'))),
|
||||
onSettled: () => {
|
||||
invalidateRowCount(queryClient, tableId)
|
||||
},
|
||||
|
||||
@@ -10,4 +10,5 @@ export {
|
||||
resolvePlanTier,
|
||||
type UpgradeCardId,
|
||||
} from './plan-view'
|
||||
export { useLimitUpgradeToast } from './use-limit-upgrade-toast'
|
||||
export { getFilledPillColor, getSubscriptionAccessState } from './utils'
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
'use client'
|
||||
|
||||
import { useCallback } from 'react'
|
||||
import { useParams, useRouter } from 'next/navigation'
|
||||
import { toast } from '@/components/emcn'
|
||||
import { buildUpgradeHref, type UpgradeReason } from '@/lib/billing/upgrade-reasons'
|
||||
|
||||
/**
|
||||
* Returns a callback that surfaces a usage-limit error as an actionable toast
|
||||
* with an "Upgrade" button deep-linking to the reason-tagged upgrade page.
|
||||
*
|
||||
* The toast persists until dismissed (emcn keeps actionable toasts open), so the
|
||||
* user always has the upgrade path within reach when they hit a limit.
|
||||
*/
|
||||
export function useLimitUpgradeToast() {
|
||||
const router = useRouter()
|
||||
const { workspaceId } = useParams<{ workspaceId: string }>()
|
||||
|
||||
return useCallback(
|
||||
(reason: UpgradeReason, message: string) => {
|
||||
toast.error(message, {
|
||||
action: {
|
||||
label: 'Upgrade',
|
||||
onClick: () => router.push(buildUpgradeHref(workspaceId, reason)),
|
||||
},
|
||||
})
|
||||
},
|
||||
[router, workspaceId]
|
||||
)
|
||||
}
|
||||
@@ -0,0 +1,169 @@
|
||||
/**
|
||||
* @vitest-environment node
|
||||
*/
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
const {
|
||||
billingFlag,
|
||||
mockClaim,
|
||||
mockSelectRows,
|
||||
dbUpdateSpy,
|
||||
sendEmailSpy,
|
||||
getEmailPreferencesMock,
|
||||
renderMock,
|
||||
subjectMock,
|
||||
isOrgAdminRoleMock,
|
||||
} = vi.hoisted(() => ({
|
||||
billingFlag: { enabled: true },
|
||||
mockClaim: vi.fn<[], unknown[]>(() => [{ id: 'u1' }]),
|
||||
mockSelectRows: vi.fn<[], unknown[]>(() => []),
|
||||
dbUpdateSpy: vi.fn(),
|
||||
sendEmailSpy: vi.fn(() => Promise.resolve({ success: true })),
|
||||
getEmailPreferencesMock: vi.fn(() => Promise.resolve(null as unknown)),
|
||||
renderMock: vi.fn(() => Promise.resolve('<html></html>')),
|
||||
subjectMock: vi.fn(() => 'Subject'),
|
||||
isOrgAdminRoleMock: vi.fn(() => true),
|
||||
}))
|
||||
|
||||
vi.mock('@sim/db', () => {
|
||||
const updateBuilder: Record<string, unknown> = {
|
||||
set: () => updateBuilder,
|
||||
where: () => updateBuilder,
|
||||
returning: () => Promise.resolve(mockClaim()),
|
||||
then: (f: (v: unknown) => unknown, r?: (e: unknown) => unknown) =>
|
||||
Promise.resolve(undefined).then(f, r),
|
||||
}
|
||||
const selectBuilder: Record<string, unknown> = {
|
||||
from: () => selectBuilder,
|
||||
where: () => selectBuilder,
|
||||
innerJoin: () => selectBuilder,
|
||||
leftJoin: () => selectBuilder,
|
||||
limit: () => Promise.resolve(mockSelectRows()),
|
||||
then: (f: (v: unknown) => unknown, r?: (e: unknown) => unknown) =>
|
||||
Promise.resolve(mockSelectRows()).then(f, r),
|
||||
}
|
||||
dbUpdateSpy.mockImplementation(() => updateBuilder)
|
||||
return { db: { update: dbUpdateSpy, select: () => selectBuilder } }
|
||||
})
|
||||
|
||||
vi.mock('@/lib/core/config/env-flags', () => ({
|
||||
get isBillingEnabled() {
|
||||
return billingFlag.enabled
|
||||
},
|
||||
}))
|
||||
vi.mock('@/lib/core/utils/urls', () => ({ getBaseUrl: () => 'https://app.sim.ai' }))
|
||||
vi.mock('@/lib/messaging/email/mailer', () => ({ sendEmail: sendEmailSpy }))
|
||||
vi.mock('@/lib/messaging/email/unsubscribe', () => ({
|
||||
getEmailPreferences: getEmailPreferencesMock,
|
||||
}))
|
||||
vi.mock('@/components/emails/render', () => ({
|
||||
renderLimitThresholdEmail: renderMock,
|
||||
getLimitEmailSubject: subjectMock,
|
||||
}))
|
||||
vi.mock('@sim/platform-authz/workspace', () => ({ isOrgAdminRole: isOrgAdminRoleMock }))
|
||||
|
||||
import { maybeSendLimitThresholdEmail } from '@/lib/billing/core/limit-notifications'
|
||||
|
||||
const baseUserParams = {
|
||||
category: 'storage' as const,
|
||||
scope: 'user' as const,
|
||||
workspaceId: 'ws-1',
|
||||
usageLabel: '4.5 GB',
|
||||
limitLabel: '5 GB',
|
||||
userId: 'u1',
|
||||
userEmail: 'u1@example.com',
|
||||
userName: 'Ada',
|
||||
}
|
||||
|
||||
describe('maybeSendLimitThresholdEmail', () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
billingFlag.enabled = true
|
||||
mockClaim.mockReturnValue([{ id: 'u1' }])
|
||||
mockSelectRows.mockReturnValue([])
|
||||
getEmailPreferencesMock.mockResolvedValue(null)
|
||||
})
|
||||
|
||||
it('sends a warning email when crossing 80% and the claim wins', async () => {
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4.5, limit: 5 })
|
||||
expect(sendEmailSpy).toHaveBeenCalledTimes(1)
|
||||
expect(renderMock).toHaveBeenCalledWith(expect.objectContaining({ kind: 'warning' }))
|
||||
expect(subjectMock).toHaveBeenCalledWith('storage', 'warning')
|
||||
})
|
||||
|
||||
it('sends a reached email at/over 100%', async () => {
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 5, limit: 5 })
|
||||
expect(renderMock).toHaveBeenCalledWith(expect.objectContaining({ kind: 'reached' }))
|
||||
expect(subjectMock).toHaveBeenCalledWith('storage', 'reached')
|
||||
})
|
||||
|
||||
it('never sends in rearmOnly mode, even when usage is above a threshold', async () => {
|
||||
await maybeSendLimitThresholdEmail({
|
||||
...baseUserParams,
|
||||
currentUsage: 4.5,
|
||||
limit: 5,
|
||||
rearmOnly: true,
|
||||
})
|
||||
expect(mockClaim).not.toHaveBeenCalled()
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('does not send when the atomic claim is lost (already notified)', async () => {
|
||||
mockClaim.mockReturnValue([])
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4.5, limit: 5 })
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('claims without re-arming on a crossing (re-arm and claim are mutually exclusive)', async () => {
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4.5, limit: 5 })
|
||||
expect(dbUpdateSpy).toHaveBeenCalledTimes(1)
|
||||
expect(mockClaim).toHaveBeenCalledTimes(1)
|
||||
expect(sendEmailSpy).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('does not send in the dead band (70%–80%)', async () => {
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 3.75, limit: 5 })
|
||||
expect(mockClaim).not.toHaveBeenCalled()
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('re-arms below the band without claiming or sending', async () => {
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 1, limit: 5 })
|
||||
expect(mockClaim).not.toHaveBeenCalled()
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('does not send OR burn the claim when the per-user toggle is off', async () => {
|
||||
mockSelectRows.mockReturnValue([{ enabled: false }])
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4.5, limit: 5 })
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
expect(mockClaim).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('does not send OR burn the claim when the recipient unsubscribed', async () => {
|
||||
getEmailPreferencesMock.mockResolvedValue({ unsubscribeNotifications: true })
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4.5, limit: 5 })
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
expect(mockClaim).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('skips entirely when billing is disabled', async () => {
|
||||
billingFlag.enabled = false
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 5, limit: 5 })
|
||||
expect(mockClaim).not.toHaveBeenCalled()
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('re-arms but does not send when usage is fully cleared (zero usage)', async () => {
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 0, limit: 5 })
|
||||
expect(dbUpdateSpy).toHaveBeenCalledTimes(1)
|
||||
expect(mockClaim).not.toHaveBeenCalled()
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('skips when the limit is non-positive', async () => {
|
||||
await maybeSendLimitThresholdEmail({ ...baseUserParams, currentUsage: 4, limit: 0 })
|
||||
expect(dbUpdateSpy).not.toHaveBeenCalled()
|
||||
expect(sendEmailSpy).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,327 @@
|
||||
import { db } from '@sim/db'
|
||||
import { member, organization, settings, user, userStats } from '@sim/db/schema'
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { isOrgAdminRole } from '@sim/platform-authz/workspace'
|
||||
import { and, eq, sql } from 'drizzle-orm'
|
||||
import { getLimitEmailSubject, renderLimitThresholdEmail } from '@/components/emails/render'
|
||||
import type { HighestPrioritySubscription } from '@/lib/billing/core/plan'
|
||||
import { getHighestPrioritySubscription } from '@/lib/billing/core/subscription'
|
||||
import { isOrgScopedSubscription } from '@/lib/billing/subscriptions/utils'
|
||||
import { buildUpgradeHref, type UpgradeReason } from '@/lib/billing/upgrade-reasons'
|
||||
import { isBillingEnabled } from '@/lib/core/config/env-flags'
|
||||
import { getBaseUrl } from '@/lib/core/utils/urls'
|
||||
import { sendEmail } from '@/lib/messaging/email/mailer'
|
||||
import { getEmailPreferences } from '@/lib/messaging/email/unsubscribe'
|
||||
|
||||
const logger = createLogger('LimitNotifications')
|
||||
|
||||
/** Limit categories that send per-category threshold emails (credits has its own path). */
|
||||
export type LimitCategory = Extract<UpgradeReason, 'storage' | 'tables' | 'seats'>
|
||||
|
||||
const WARN_THRESHOLD = 80
|
||||
const REACH_THRESHOLD = 100
|
||||
/** Usage must drop below this band before the same threshold can re-notify (hysteresis). */
|
||||
const REARM_BELOW = 70
|
||||
|
||||
/**
|
||||
* Resolve the threshold a given usage percent should be notified at:
|
||||
* 100 at/over the limit, 80 when approaching, 0 otherwise.
|
||||
*/
|
||||
function thresholdFor(percent: number): 0 | 80 | 100 {
|
||||
if (percent >= REACH_THRESHOLD) return REACH_THRESHOLD
|
||||
if (percent >= WARN_THRESHOLD) return WARN_THRESHOLD
|
||||
return 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Atomically claim a threshold for a category: advance the stored value to
|
||||
* `threshold` only if it is currently lower, returning whether THIS call won the
|
||||
* advance. A single conditional UPDATE is race-free — concurrent crossings can't
|
||||
* both claim, so the email is sent exactly once per crossing.
|
||||
*/
|
||||
async function claimThreshold(
|
||||
scope: 'user' | 'organization',
|
||||
id: string,
|
||||
category: LimitCategory,
|
||||
threshold: number
|
||||
): Promise<boolean> {
|
||||
const setExpr = sql`jsonb_set(coalesce(${scope === 'user' ? userStats.limitNotifications : organization.limitNotifications}, '{}'::jsonb), ARRAY[${category}], to_jsonb(${threshold}::int))`
|
||||
const onlyIfLower =
|
||||
scope === 'user'
|
||||
? sql`coalesce((${userStats.limitNotifications} ->> ${category})::int, 0) < ${threshold}`
|
||||
: sql`coalesce((${organization.limitNotifications} ->> ${category})::int, 0) < ${threshold}`
|
||||
|
||||
const claimed =
|
||||
scope === 'user'
|
||||
? await db
|
||||
.update(userStats)
|
||||
.set({ limitNotifications: setExpr })
|
||||
.where(and(eq(userStats.userId, id), onlyIfLower))
|
||||
.returning({ id: userStats.userId })
|
||||
: await db
|
||||
.update(organization)
|
||||
.set({ limitNotifications: setExpr })
|
||||
.where(and(eq(organization.id, id), onlyIfLower))
|
||||
.returning({ id: organization.id })
|
||||
|
||||
return claimed.length > 0
|
||||
}
|
||||
|
||||
/** Re-arm a category (reset its stored threshold to 0) once usage falls back into the low band. */
|
||||
async function rearmThreshold(
|
||||
scope: 'user' | 'organization',
|
||||
id: string,
|
||||
category: LimitCategory
|
||||
): Promise<void> {
|
||||
const setExpr = sql`jsonb_set(coalesce(${scope === 'user' ? userStats.limitNotifications : organization.limitNotifications}, '{}'::jsonb), ARRAY[${category}], to_jsonb(0::int))`
|
||||
const onlyIfArmed =
|
||||
scope === 'user'
|
||||
? sql`coalesce((${userStats.limitNotifications} ->> ${category})::int, 0) > 0`
|
||||
: sql`coalesce((${organization.limitNotifications} ->> ${category})::int, 0) > 0`
|
||||
|
||||
if (scope === 'user') {
|
||||
await db
|
||||
.update(userStats)
|
||||
.set({ limitNotifications: setExpr })
|
||||
.where(and(eq(userStats.userId, id), onlyIfArmed))
|
||||
} else {
|
||||
await db
|
||||
.update(organization)
|
||||
.set({ limitNotifications: setExpr })
|
||||
.where(and(eq(organization.id, id), onlyIfArmed))
|
||||
}
|
||||
}
|
||||
|
||||
interface LimitEmailRecipient {
|
||||
email: string
|
||||
name?: string
|
||||
}
|
||||
|
||||
/** Whether a recipient has unsubscribed from all or from notification emails. */
|
||||
async function isUnsubscribed(email: string): Promise<boolean> {
|
||||
const prefs = await getEmailPreferences(email)
|
||||
return Boolean(prefs?.unsubscribeAll || prefs?.unsubscribeNotifications)
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve the recipients that should receive a limit email, with all opt-outs
|
||||
* already applied (the per-user notifications toggle and unsubscribe prefs).
|
||||
* Returning an empty list means "nobody to notify" — the caller then skips the
|
||||
* claim so the dedup state isn't burned without an email going out.
|
||||
*/
|
||||
async function resolveRecipients(
|
||||
scope: 'user' | 'organization',
|
||||
params: { userId?: string; userEmail?: string; userName?: string; organizationId?: string }
|
||||
): Promise<LimitEmailRecipient[]> {
|
||||
if (scope === 'user') {
|
||||
if (!params.userId || !params.userEmail) return []
|
||||
const rows = await db
|
||||
.select({ enabled: settings.billingUsageNotificationsEnabled })
|
||||
.from(settings)
|
||||
.where(eq(settings.userId, params.userId))
|
||||
.limit(1)
|
||||
if (rows.length > 0 && rows[0].enabled === false) return []
|
||||
if (await isUnsubscribed(params.userEmail)) return []
|
||||
return [{ email: params.userEmail, name: params.userName }]
|
||||
}
|
||||
|
||||
if (!params.organizationId) return []
|
||||
const admins = await db
|
||||
.select({
|
||||
email: user.email,
|
||||
name: user.name,
|
||||
enabled: settings.billingUsageNotificationsEnabled,
|
||||
role: member.role,
|
||||
})
|
||||
.from(member)
|
||||
.innerJoin(user, eq(member.userId, user.id))
|
||||
.leftJoin(settings, eq(settings.userId, member.userId))
|
||||
.where(eq(member.organizationId, params.organizationId))
|
||||
|
||||
const recipients: LimitEmailRecipient[] = []
|
||||
for (const a of admins) {
|
||||
if (!isOrgAdminRole(a.role)) continue
|
||||
if (a.enabled === false) continue
|
||||
if (!a.email) continue
|
||||
if (await isUnsubscribed(a.email)) continue
|
||||
recipients.push({ email: a.email, name: a.name || undefined })
|
||||
}
|
||||
return recipients
|
||||
}
|
||||
|
||||
/**
|
||||
* Send a usage-limit threshold email (80% warning / 100% reached) for a
|
||||
* non-credit category, edge-triggered on the mutation that changed usage.
|
||||
*
|
||||
* Flow: bail when billing is off or the limit is non-positive; re-arm the
|
||||
* persisted threshold when current usage is back in the low band; then (for
|
||||
* increases) resolve eligible recipients and atomically claim the threshold
|
||||
* before sending. Re-arm and claim are mutually exclusive per call — re-arm only
|
||||
* fires when `desired === 0` — so the dedup stays a single atomic
|
||||
* {@link claimThreshold} with no re-arm/claim interleaving race. Recipients are
|
||||
* resolved with opt-outs applied BEFORE the claim, so an opted-out recipient
|
||||
* never burns the threshold (which would suppress a later email once
|
||||
* notifications are re-enabled). Per-recipient send failures are isolated.
|
||||
*
|
||||
* The highest threshold already emailed is persisted per category on
|
||||
* `user_stats` / `organization`; it re-arms once usage drops below
|
||||
* {@link REARM_BELOW}. Best-effort — callers fire-and-forget; failures never
|
||||
* block the mutation. Mirrors the credits path in `maybeSendUsageThresholdEmail`:
|
||||
* respects the per-user notifications toggle and unsubscribe preferences, and
|
||||
* emails org admins for organization-scoped limits.
|
||||
*/
|
||||
export async function maybeSendLimitThresholdEmail(params: {
|
||||
category: LimitCategory
|
||||
scope: 'user' | 'organization'
|
||||
workspaceId: string
|
||||
currentUsage: number
|
||||
limit: number
|
||||
/** Pre-formatted current usage for the email body, e.g. "4.2 GB", "9 seats". */
|
||||
usageLabel: string
|
||||
/** Pre-formatted limit for the email body, e.g. "5 GB", "10 seats". */
|
||||
limitLabel: string
|
||||
/**
|
||||
* When true, only the re-arm is evaluated and no email is ever sent. Used by
|
||||
* usage-decrease paths (e.g. a storage shrink) where usage can still be above
|
||||
* a threshold but the change is a drop, not a fresh crossing.
|
||||
*/
|
||||
rearmOnly?: boolean
|
||||
userId?: string
|
||||
userEmail?: string
|
||||
userName?: string
|
||||
organizationId?: string
|
||||
}): Promise<void> {
|
||||
try {
|
||||
if (!isBillingEnabled) return
|
||||
if (params.limit <= 0) return
|
||||
|
||||
const { category, scope } = params
|
||||
const percent = Math.max(0, (params.currentUsage / params.limit) * 100)
|
||||
const desired = thresholdFor(percent)
|
||||
|
||||
const stateId = scope === 'user' ? params.userId : params.organizationId
|
||||
if (!stateId) return
|
||||
|
||||
if (percent < REARM_BELOW) {
|
||||
await rearmThreshold(scope, stateId, category)
|
||||
}
|
||||
|
||||
if (params.rearmOnly || desired === 0) return
|
||||
|
||||
const recipients = await resolveRecipients(scope, params)
|
||||
if (recipients.length === 0) return
|
||||
|
||||
if (!(await claimThreshold(scope, stateId, category, desired))) return
|
||||
|
||||
const kind = desired === REACH_THRESHOLD ? 'reached' : 'warning'
|
||||
const percentUsed = Math.min(100, Math.round(percent))
|
||||
const upgradeLink = `${getBaseUrl()}${buildUpgradeHref(params.workspaceId, category)}`
|
||||
|
||||
let sent = 0
|
||||
for (const r of recipients) {
|
||||
try {
|
||||
const html = await renderLimitThresholdEmail({
|
||||
kind,
|
||||
reason: category,
|
||||
userName: r.name,
|
||||
usageLabel: params.usageLabel,
|
||||
limitLabel: params.limitLabel,
|
||||
percentUsed,
|
||||
upgradeLink,
|
||||
})
|
||||
|
||||
await sendEmail({
|
||||
to: r.email,
|
||||
subject: getLimitEmailSubject(category, kind),
|
||||
html,
|
||||
emailType: 'notifications',
|
||||
})
|
||||
sent++
|
||||
} catch (sendError) {
|
||||
logger.error('Failed to send limit email', {
|
||||
category,
|
||||
email: r.email,
|
||||
error: sendError,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
if (sent > 0) {
|
||||
logger.info('Sent usage-limit threshold email', { category, scope, kind, percentUsed, sent })
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error('Failed to send usage-limit threshold email', {
|
||||
category: params.category,
|
||||
scope: params.scope,
|
||||
error,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve billing scope for a billed account user, then dispatch the limit
|
||||
* threshold email via {@link maybeSendLimitThresholdEmail}.
|
||||
*
|
||||
* The single entry point for per-category usage-limit emails: callers supply the
|
||||
* billed user, the usage numbers, and pre-formatted labels, and this resolves
|
||||
* personal vs. pooled (org) scope and the recipient. Best-effort — never throws.
|
||||
*
|
||||
* @param params.billedUserId - User whose subscription determines billing scope
|
||||
* (the uploader for storage; the workspace's billed account for tables).
|
||||
* @param params.subscription - Pre-resolved subscription (may be `null`) to skip
|
||||
* the `getHighestPrioritySubscription` lookup on hot paths; omit to fetch here.
|
||||
*/
|
||||
export async function maybeNotifyLimit(params: {
|
||||
category: LimitCategory
|
||||
billedUserId: string
|
||||
workspaceId: string
|
||||
currentUsage: number
|
||||
limit: number
|
||||
usageLabel: string
|
||||
limitLabel: string
|
||||
/** Re-arm only, never send — for usage-decrease callers. See {@link maybeSendLimitThresholdEmail}. */
|
||||
rearmOnly?: boolean
|
||||
subscription?: HighestPrioritySubscription | null
|
||||
}): Promise<void> {
|
||||
try {
|
||||
const sub =
|
||||
params.subscription === undefined
|
||||
? await getHighestPrioritySubscription(params.billedUserId)
|
||||
: params.subscription
|
||||
const isOrg = Boolean(sub && isOrgScopedSubscription(sub, params.billedUserId))
|
||||
|
||||
const percent = params.limit > 0 ? (params.currentUsage / params.limit) * 100 : 0
|
||||
let userEmail: string | undefined
|
||||
let userName: string | undefined
|
||||
if (!isOrg && !params.rearmOnly && percent >= WARN_THRESHOLD) {
|
||||
const [row] = await db
|
||||
.select({ email: user.email, name: user.name })
|
||||
.from(user)
|
||||
.where(eq(user.id, params.billedUserId))
|
||||
.limit(1)
|
||||
userEmail = row?.email
|
||||
userName = row?.name || undefined
|
||||
}
|
||||
|
||||
await maybeSendLimitThresholdEmail({
|
||||
category: params.category,
|
||||
scope: isOrg ? 'organization' : 'user',
|
||||
workspaceId: params.workspaceId,
|
||||
currentUsage: params.currentUsage,
|
||||
limit: params.limit,
|
||||
usageLabel: params.usageLabel,
|
||||
limitLabel: params.limitLabel,
|
||||
rearmOnly: params.rearmOnly,
|
||||
userId: params.billedUserId,
|
||||
userEmail,
|
||||
userName,
|
||||
organizationId: isOrg ? sub?.referenceId : undefined,
|
||||
})
|
||||
} catch (error) {
|
||||
logger.error('Failed to resolve scope for usage-limit notification', {
|
||||
category: params.category,
|
||||
billedUserId: params.billedUserId,
|
||||
error,
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -32,6 +32,7 @@ import {
|
||||
isOrgScopedSubscription,
|
||||
} from '@/lib/billing/subscriptions/utils'
|
||||
import type { BillingData, UsageData, UsageLimitInfo } from '@/lib/billing/types'
|
||||
import { buildUpgradeHref } from '@/lib/billing/upgrade-reasons'
|
||||
import { Decimal, toDecimal, toNumber } from '@/lib/billing/utils/decimal'
|
||||
import { isBillingEnabled } from '@/lib/core/config/env-flags'
|
||||
import { getBaseUrl } from '@/lib/core/utils/urls'
|
||||
@@ -828,6 +829,8 @@ export async function maybeSendUsageThresholdEmail(params: {
|
||||
userEmail?: string
|
||||
userName?: string
|
||||
organizationId?: string
|
||||
/** Workspace the usage occurred in, used to build a live upgrade/billing link. */
|
||||
workspaceId?: string
|
||||
currentUsageAfter: number
|
||||
limit: number
|
||||
}): Promise<void> {
|
||||
@@ -838,6 +841,13 @@ export async function maybeSendUsageThresholdEmail(params: {
|
||||
const baseUrl = getBaseUrl()
|
||||
const isFreeUser = params.planName === 'Free'
|
||||
|
||||
const upgradeCreditsLink = params.workspaceId
|
||||
? `${baseUrl}${buildUpgradeHref(params.workspaceId, 'credits')}`
|
||||
: `${baseUrl}/workspace`
|
||||
const billingSettingsLink = params.workspaceId
|
||||
? `${baseUrl}/workspace/${params.workspaceId}/settings/billing`
|
||||
: `${baseUrl}/workspace`
|
||||
|
||||
// Check for 80% threshold crossing — used for paid users (budget warning) and free users (upgrade nudge)
|
||||
const crosses80 = params.percentBefore < 80 && params.percentAfter >= 80
|
||||
// Check for 100% threshold (free users only — credits exhausted)
|
||||
@@ -848,7 +858,7 @@ export async function maybeSendUsageThresholdEmail(params: {
|
||||
|
||||
// For 80% threshold email (paid users only)
|
||||
if (crosses80 && !isFreeUser) {
|
||||
const ctaLink = `${baseUrl}/workspace?billing=usage`
|
||||
const ctaLink = billingSettingsLink
|
||||
const sendTo = async (email: string, name?: string) => {
|
||||
const prefs = await getEmailPreferences(email)
|
||||
if (prefs?.unsubscribeAll || prefs?.unsubscribeNotifications) return
|
||||
@@ -903,7 +913,7 @@ export async function maybeSendUsageThresholdEmail(params: {
|
||||
|
||||
// For 80% threshold email (free users only — skip if they also crossed 100% in same call)
|
||||
if (crosses80 && isFreeUser && !crosses100) {
|
||||
const upgradeLink = `${baseUrl}/workspace?billing=upgrade`
|
||||
const upgradeLink = upgradeCreditsLink
|
||||
const sendFreeTierEmail = async (email: string, name?: string) => {
|
||||
const prefs = await getEmailPreferences(email)
|
||||
if (prefs?.unsubscribeAll || prefs?.unsubscribeNotifications) return
|
||||
@@ -945,7 +955,7 @@ export async function maybeSendUsageThresholdEmail(params: {
|
||||
|
||||
// For 100% threshold email (free users only — credits exhausted)
|
||||
if (crosses100 && isFreeUser) {
|
||||
const upgradeLink = `${baseUrl}/workspace?billing=upgrade`
|
||||
const upgradeLink = upgradeCreditsLink
|
||||
const sendExhaustedEmail = async (email: string, name?: string) => {
|
||||
const prefs = await getEmailPreferences(email)
|
||||
if (prefs?.unsubscribeAll || prefs?.unsubscribeNotifications) return
|
||||
|
||||
@@ -10,9 +10,10 @@ import {
|
||||
DEFAULT_PRO_STORAGE_LIMIT_GB,
|
||||
DEFAULT_TEAM_STORAGE_LIMIT_GB,
|
||||
} from '@sim/db/constants'
|
||||
import { organization, subscription, userStats } from '@sim/db/schema'
|
||||
import { organization, userStats } from '@sim/db/schema'
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { eq } from 'drizzle-orm'
|
||||
import type { HighestPrioritySubscription } from '@/lib/billing/core/plan'
|
||||
import { getPlanTypeForLimits, isEnterprise, isFree } from '@/lib/billing/plan-helpers'
|
||||
import { isOrgScopedSubscription } from '@/lib/billing/subscriptions/utils'
|
||||
import { getEnv } from '@/lib/core/config/env'
|
||||
@@ -20,6 +21,12 @@ import { isBillingEnabled } from '@/lib/core/config/env-flags'
|
||||
|
||||
const logger = createLogger('StorageLimits')
|
||||
|
||||
/** Resolve the highest-priority subscription via a deferred import (avoids a static cycle). */
|
||||
async function resolveSub(userId: string): Promise<HighestPrioritySubscription | null> {
|
||||
const { getHighestPrioritySubscription } = await import('@/lib/billing/core/subscription')
|
||||
return getHighestPrioritySubscription(userId)
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert GB to bytes
|
||||
*/
|
||||
@@ -74,13 +81,19 @@ export function getStorageLimitForPlan(plan: string, metadata?: any): number {
|
||||
}
|
||||
|
||||
/**
|
||||
* Get storage limit for a user based on their subscription
|
||||
* Returns limit in bytes
|
||||
* Get storage limit for a user based on their subscription. Returns limit in
|
||||
* bytes.
|
||||
*
|
||||
* @param prefetchedSub - Pass an already-resolved subscription (may be `null`)
|
||||
* to skip the `getHighestPrioritySubscription` lookup on hot paths. Omit
|
||||
* (leave `undefined`) to fetch it here.
|
||||
*/
|
||||
export async function getUserStorageLimit(userId: string): Promise<number> {
|
||||
export async function getUserStorageLimit(
|
||||
userId: string,
|
||||
prefetchedSub?: HighestPrioritySubscription | null
|
||||
): Promise<number> {
|
||||
try {
|
||||
const { getHighestPrioritySubscription } = await import('@/lib/billing/core/subscription')
|
||||
const sub = await getHighestPrioritySubscription(userId)
|
||||
const sub = prefetchedSub === undefined ? await resolveSub(userId) : prefetchedSub
|
||||
|
||||
const limits = getStorageLimits()
|
||||
|
||||
@@ -88,22 +101,11 @@ export async function getUserStorageLimit(userId: string): Promise<number> {
|
||||
return limits.free
|
||||
}
|
||||
|
||||
// Org-scoped subs use pooled org-level storage. Custom limits come from the
|
||||
// subscription metadata; otherwise use the team/enterprise default.
|
||||
if (isOrgScopedSubscription(sub, userId)) {
|
||||
const orgRecord = await db
|
||||
.select({ metadata: subscription.metadata })
|
||||
.from(subscription)
|
||||
.where(eq(subscription.id, sub.id))
|
||||
.limit(1)
|
||||
|
||||
if (orgRecord.length > 0 && orgRecord[0].metadata) {
|
||||
const metadata = orgRecord[0].metadata as any
|
||||
if (metadata.customStorageLimitGB) {
|
||||
return metadata.customStorageLimitGB * 1024 * 1024 * 1024
|
||||
}
|
||||
const metadata = sub.metadata as { customStorageLimitGB?: number } | null
|
||||
if (metadata?.customStorageLimitGB) {
|
||||
return metadata.customStorageLimitGB * 1024 * 1024 * 1024
|
||||
}
|
||||
|
||||
return isEnterprise(sub.plan) ? limits.enterpriseDefault : limits.team
|
||||
}
|
||||
|
||||
@@ -122,13 +124,17 @@ export async function getUserStorageLimit(userId: string): Promise<number> {
|
||||
}
|
||||
|
||||
/**
|
||||
* Get current storage usage for a user
|
||||
* Returns usage in bytes
|
||||
* Get current storage usage for a user. Returns usage in bytes.
|
||||
*
|
||||
* @param prefetchedSub - Pass an already-resolved subscription (may be `null`)
|
||||
* to skip the `getHighestPrioritySubscription` lookup on hot paths.
|
||||
*/
|
||||
export async function getUserStorageUsage(userId: string): Promise<number> {
|
||||
export async function getUserStorageUsage(
|
||||
userId: string,
|
||||
prefetchedSub?: HighestPrioritySubscription | null
|
||||
): Promise<number> {
|
||||
try {
|
||||
const { getHighestPrioritySubscription } = await import('@/lib/billing/core/subscription')
|
||||
const sub = await getHighestPrioritySubscription(userId)
|
||||
const sub = prefetchedSub === undefined ? await resolveSub(userId) : prefetchedSub
|
||||
|
||||
// Org-scoped subs share pooled `organization.storageUsedBytes`;
|
||||
// personal plans use `userStats`.
|
||||
|
||||
@@ -8,24 +8,80 @@ import { db } from '@sim/db'
|
||||
import { organization, userStats } from '@sim/db/schema'
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { eq, sql } from 'drizzle-orm'
|
||||
import { maybeNotifyLimit } from '@/lib/billing/core/limit-notifications'
|
||||
import type { HighestPrioritySubscription } from '@/lib/billing/core/plan'
|
||||
import { getUserStorageLimit, getUserStorageUsage } from '@/lib/billing/storage/limits'
|
||||
import { isOrgScopedSubscription } from '@/lib/billing/subscriptions/utils'
|
||||
import { isBillingEnabled } from '@/lib/core/config/env-flags'
|
||||
|
||||
const logger = createLogger('StorageTracking')
|
||||
|
||||
/** Format bytes as a `GB` label for usage-limit emails (2dp usage, whole-number limit). */
|
||||
function formatGb(bytes: number, decimals: number): string {
|
||||
return `${(bytes / 1024 ** 3).toFixed(decimals)} GB`
|
||||
}
|
||||
|
||||
/**
|
||||
* Best-effort storage threshold evaluation after a usage change. Re-reads the
|
||||
* (now updated) usage and plan limit, then delegates dedup + send to
|
||||
* {@link maybeNotifyLimit}. Never throws.
|
||||
*
|
||||
* The caller passes the subscription it already resolved for the increment/
|
||||
* decrement, so the whole path (usage read, limit, scope) reuses a single
|
||||
* `getHighestPrioritySubscription` instead of re-fetching it three times.
|
||||
*
|
||||
* @param rearmOnly - True on decrements, so a shrink that leaves usage above a
|
||||
* threshold re-arms but never sends (a drop is not a fresh crossing).
|
||||
*/
|
||||
async function maybeNotifyStorageLimit(
|
||||
userId: string,
|
||||
workspaceId: string,
|
||||
sub: HighestPrioritySubscription | null,
|
||||
rearmOnly = false
|
||||
): Promise<void> {
|
||||
try {
|
||||
const [usage, limit] = await Promise.all([
|
||||
getUserStorageUsage(userId, sub),
|
||||
getUserStorageLimit(userId, sub),
|
||||
])
|
||||
|
||||
await maybeNotifyLimit({
|
||||
category: 'storage',
|
||||
billedUserId: userId,
|
||||
workspaceId,
|
||||
currentUsage: usage,
|
||||
limit,
|
||||
usageLabel: formatGb(usage, 2),
|
||||
limitLabel: formatGb(limit, 0),
|
||||
rearmOnly,
|
||||
subscription: sub,
|
||||
})
|
||||
} catch (error) {
|
||||
logger.error('Error evaluating storage limit notification:', error)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Increment storage usage after successful file upload
|
||||
* Only tracks if billing is enabled
|
||||
*
|
||||
* @param workspaceId - When provided, evaluates the storage usage-limit email
|
||||
* (80% / 100%) after the increment. Best-effort; never blocks the upload.
|
||||
*/
|
||||
export async function incrementStorageUsage(userId: string, bytes: number): Promise<void> {
|
||||
export async function incrementStorageUsage(
|
||||
userId: string,
|
||||
bytes: number,
|
||||
workspaceId?: string
|
||||
): Promise<void> {
|
||||
if (!isBillingEnabled) {
|
||||
logger.debug('Billing disabled, skipping storage increment')
|
||||
return
|
||||
}
|
||||
|
||||
let sub: HighestPrioritySubscription | null = null
|
||||
try {
|
||||
const { getHighestPrioritySubscription } = await import('@/lib/billing/core/subscription')
|
||||
const sub = await getHighestPrioritySubscription(userId)
|
||||
sub = await getHighestPrioritySubscription(userId)
|
||||
|
||||
// Org-scoped subs pool at the org level; personal plans per-user.
|
||||
if (isOrgScopedSubscription(sub, userId) && sub) {
|
||||
@@ -51,21 +107,35 @@ export async function incrementStorageUsage(userId: string, bytes: number): Prom
|
||||
logger.error('Error incrementing storage usage:', error)
|
||||
throw error
|
||||
}
|
||||
|
||||
if (workspaceId) {
|
||||
void maybeNotifyStorageLimit(userId, workspaceId, sub)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Decrement storage usage after file deletion
|
||||
* Only tracks if billing is enabled
|
||||
*
|
||||
* @param workspaceId - When provided, re-evaluates the storage threshold state
|
||||
* after the decrement. Usage only drops here, so this can only re-arm a
|
||||
* previously-sent threshold (it never sends), keeping the re-warning correct
|
||||
* after a shrink. Best-effort; never blocks the caller.
|
||||
*/
|
||||
export async function decrementStorageUsage(userId: string, bytes: number): Promise<void> {
|
||||
export async function decrementStorageUsage(
|
||||
userId: string,
|
||||
bytes: number,
|
||||
workspaceId?: string
|
||||
): Promise<void> {
|
||||
if (!isBillingEnabled) {
|
||||
logger.debug('Billing disabled, skipping storage decrement')
|
||||
return
|
||||
}
|
||||
|
||||
let sub: HighestPrioritySubscription | null = null
|
||||
try {
|
||||
const { getHighestPrioritySubscription } = await import('@/lib/billing/core/subscription')
|
||||
const sub = await getHighestPrioritySubscription(userId)
|
||||
sub = await getHighestPrioritySubscription(userId)
|
||||
|
||||
if (isOrgScopedSubscription(sub, userId) && sub) {
|
||||
await db
|
||||
@@ -90,4 +160,8 @@ export async function decrementStorageUsage(userId: string, bytes: number): Prom
|
||||
logger.error('Error decrementing storage usage:', error)
|
||||
throw error
|
||||
}
|
||||
|
||||
if (workspaceId) {
|
||||
void maybeNotifyStorageLimit(userId, workspaceId, sub, true)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
/**
|
||||
* @vitest-environment node
|
||||
*/
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import {
|
||||
buildUpgradeHref,
|
||||
isUpgradeReason,
|
||||
UPGRADE_REASON_COPY,
|
||||
UPGRADE_REASONS,
|
||||
} from '@/lib/billing/upgrade-reasons'
|
||||
|
||||
describe('upgrade-reasons', () => {
|
||||
it('has copy for every reason', () => {
|
||||
for (const reason of UPGRADE_REASONS) {
|
||||
const copy = UPGRADE_REASON_COPY[reason]
|
||||
expect(copy.header).toMatch(/^Upgrade to scale/)
|
||||
expect(copy.noun.length).toBeGreaterThan(0)
|
||||
expect(copy.warningSubject.length).toBeGreaterThan(0)
|
||||
expect(copy.reachedSubject.length).toBeGreaterThan(0)
|
||||
}
|
||||
})
|
||||
|
||||
it('uses Emir’s header wording', () => {
|
||||
expect(UPGRADE_REASON_COPY.seats.header).toBe('Upgrade to scale with your teammates')
|
||||
expect(UPGRADE_REASON_COPY.tables.header).toBe('Upgrade to scale your tables')
|
||||
expect(UPGRADE_REASON_COPY.storage.header).toBe('Upgrade to scale your storage')
|
||||
})
|
||||
|
||||
it('builds hrefs with and without a reason', () => {
|
||||
expect(buildUpgradeHref('ws-1')).toBe('/workspace/ws-1/upgrade')
|
||||
expect(buildUpgradeHref('ws-1', 'tables')).toBe('/workspace/ws-1/upgrade?reason=tables')
|
||||
})
|
||||
|
||||
it('guards known reasons', () => {
|
||||
expect(isUpgradeReason('storage')).toBe(true)
|
||||
expect(isUpgradeReason('seats')).toBe(true)
|
||||
expect(isUpgradeReason('bogus')).toBe(false)
|
||||
expect(isUpgradeReason(null)).toBe(false)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,88 @@
|
||||
/**
|
||||
* Upgrade-reason registry.
|
||||
*
|
||||
* Single source of truth for the language shown when a user is routed to the
|
||||
* upgrade page after hitting a usage limit. The same copy drives both the
|
||||
* upgrade-page header and the threshold/limit emails, so the in-app and email
|
||||
* journeys never drift apart.
|
||||
*/
|
||||
|
||||
/** The limit categories that can route a user to the upgrade page. */
|
||||
export const UPGRADE_REASONS = ['credits', 'storage', 'tables', 'seats'] as const
|
||||
|
||||
export type UpgradeReason = (typeof UPGRADE_REASONS)[number]
|
||||
|
||||
/** URL query key the upgrade page reads to resolve the reason. */
|
||||
export const UPGRADE_REASON_PARAM = 'reason' as const
|
||||
|
||||
/** Header shown on the upgrade page when no (or an invalid) reason is present. */
|
||||
export const DEFAULT_UPGRADE_HEADER = 'Plans that scale with you' as const
|
||||
|
||||
interface UpgradeReasonCopy {
|
||||
/** Upgrade-page `<h1>` header. */
|
||||
header: string
|
||||
/** Lowercase noun for the limited resource (e.g. "tables", "storage"). */
|
||||
noun: string
|
||||
/** Subject line for the 80% warning email. */
|
||||
warningSubject: string
|
||||
/** Subject line for the 100% limit-reached email. */
|
||||
reachedSubject: string
|
||||
/** One-line body lead for the warning email (running low). */
|
||||
warningLead: string
|
||||
/** One-line body lead for the limit-reached email. */
|
||||
reachedLead: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Per-reason copy. Headers follow the "Upgrade to scale ..." pattern; email
|
||||
* subjects/leads reuse the same noun so a user sees consistent language whether
|
||||
* they arrive from the app or from an email.
|
||||
*/
|
||||
export const UPGRADE_REASON_COPY: Record<UpgradeReason, UpgradeReasonCopy> = {
|
||||
credits: {
|
||||
header: 'Upgrade to scale your usage',
|
||||
noun: 'credits',
|
||||
warningSubject: "You're nearing your usage limit",
|
||||
reachedSubject: "You've reached your usage limit",
|
||||
warningLead: "You're approaching your usage limit.",
|
||||
reachedLead: "You've reached your usage limit.",
|
||||
},
|
||||
storage: {
|
||||
header: 'Upgrade to scale your storage',
|
||||
noun: 'storage',
|
||||
warningSubject: "You're running low on storage",
|
||||
reachedSubject: "You've reached your storage limit",
|
||||
warningLead: "You're running low on storage.",
|
||||
reachedLead: "You've reached your storage limit.",
|
||||
},
|
||||
tables: {
|
||||
header: 'Upgrade to scale your tables',
|
||||
noun: 'table rows',
|
||||
warningSubject: "You're running low on table space",
|
||||
reachedSubject: "You've reached your table limit",
|
||||
warningLead: "You're running low on table space.",
|
||||
reachedLead: "You've reached your table limit.",
|
||||
},
|
||||
seats: {
|
||||
header: 'Upgrade to scale with your teammates',
|
||||
noun: 'seats',
|
||||
warningSubject: "You're running low on seats",
|
||||
reachedSubject: "You've used all your seats",
|
||||
warningLead: "You're running low on seats for your team.",
|
||||
reachedLead: "You've used all the seats on your plan.",
|
||||
},
|
||||
}
|
||||
|
||||
/** Type guard for a raw query value against the known reasons. */
|
||||
export function isUpgradeReason(value: string | null | undefined): value is UpgradeReason {
|
||||
return value != null && (UPGRADE_REASONS as readonly string[]).includes(value)
|
||||
}
|
||||
|
||||
/**
|
||||
* Build a link to the workspace upgrade page, optionally tagged with the reason
|
||||
* that sent the user there so the page can swap its header.
|
||||
*/
|
||||
export function buildUpgradeHref(workspaceId: string, reason?: UpgradeReason): string {
|
||||
const base = `/workspace/${workspaceId}/upgrade`
|
||||
return reason ? `${base}?${UPGRADE_REASON_PARAM}=${reason}` : base
|
||||
}
|
||||
@@ -1006,6 +1006,7 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
userEmail: emailContext.userEmail,
|
||||
userName: emailContext.userName || undefined,
|
||||
planName: emailContext.planName,
|
||||
workspaceId: updatedLog.workspaceId,
|
||||
percentBefore,
|
||||
percentAfter,
|
||||
currentUsageAfter,
|
||||
@@ -1022,6 +1023,7 @@ export class ExecutionLogger implements IExecutionLoggerService {
|
||||
scope: 'organization',
|
||||
organizationId: emailContext.organizationId,
|
||||
planName: emailContext.planName,
|
||||
workspaceId: updatedLog.workspaceId,
|
||||
percentBefore,
|
||||
percentAfter,
|
||||
currentUsageAfter,
|
||||
|
||||
@@ -19,7 +19,8 @@ vi.mock('@sim/db', () => dbChainMock)
|
||||
// Capacity is exercised in billing.test.ts; here it's a no-op so the timeout-scaling
|
||||
// suites can use large synthetic row counts without tripping the plan limit.
|
||||
vi.mock('@/lib/table/billing', () => ({
|
||||
assertRowCapacity: vi.fn().mockResolvedValue(undefined),
|
||||
assertRowCapacity: vi.fn().mockResolvedValue(1_000_000),
|
||||
notifyTableRowUsage: vi.fn(),
|
||||
getMaxRowsPerTable: vi.fn().mockResolvedValue(1_000_000),
|
||||
wouldExceedRowLimit: () => false,
|
||||
TableRowLimitError: class TableRowLimitError extends Error {},
|
||||
|
||||
@@ -32,6 +32,7 @@ import {
|
||||
assertRowCapacity,
|
||||
getMaxRowsPerTable,
|
||||
getWorkspaceTableLimits,
|
||||
notifyTableRowUsage,
|
||||
TableRowLimitError,
|
||||
wouldExceedRowLimit,
|
||||
} from '@/lib/table/billing'
|
||||
@@ -120,16 +121,16 @@ describe('wouldExceedRowLimit', () => {
|
||||
})
|
||||
|
||||
describe('assertRowCapacity', () => {
|
||||
it('passes when the write stays under the plan limit', async () => {
|
||||
it('returns the resolved limit when the write stays under it', async () => {
|
||||
await expect(
|
||||
assertRowCapacity({ workspaceId: nextWorkspaceId(), currentRowCount: 10, addedRows: 5 })
|
||||
).resolves.toBeUndefined()
|
||||
).resolves.toBe(5000)
|
||||
})
|
||||
|
||||
it('allows reaching the limit exactly', async () => {
|
||||
it('allows reaching the limit exactly and returns it', async () => {
|
||||
await expect(
|
||||
assertRowCapacity({ workspaceId: nextWorkspaceId(), currentRowCount: 4999, addedRows: 1 })
|
||||
).resolves.toBeUndefined()
|
||||
).resolves.toBe(5000)
|
||||
})
|
||||
|
||||
it('throws TableRowLimitError when the write would exceed the limit', async () => {
|
||||
@@ -155,6 +156,35 @@ describe('assertRowCapacity', () => {
|
||||
currentRowCount: 10_000_000,
|
||||
addedRows: 1,
|
||||
})
|
||||
).resolves.toBeUndefined()
|
||||
).resolves.toBe(-1)
|
||||
})
|
||||
})
|
||||
|
||||
describe('notifyTableRowUsage — edge-crossing gate', () => {
|
||||
beforeEach(() => mockGetWorkspaceBilledAccountUserId.mockClear())
|
||||
|
||||
it('fires when an insert crosses UP into the warn band (limit 5000)', () => {
|
||||
notifyTableRowUsage({ workspaceId: 'ws', currentRowCount: 3990, addedRows: 20, limit: 5000 })
|
||||
expect(mockGetWorkspaceBilledAccountUserId).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('fires when an insert crosses UP into the reached band', () => {
|
||||
notifyTableRowUsage({ workspaceId: 'ws', currentRowCount: 4990, addedRows: 20, limit: 5000 })
|
||||
expect(mockGetWorkspaceBilledAccountUserId).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('does NOT fire when already above the band (no crossing)', () => {
|
||||
notifyTableRowUsage({ workspaceId: 'ws', currentRowCount: 4200, addedRows: 100, limit: 5000 })
|
||||
expect(mockGetWorkspaceBilledAccountUserId).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('does NOT fire well below the band', () => {
|
||||
notifyTableRowUsage({ workspaceId: 'ws', currentRowCount: 100, addedRows: 10, limit: 5000 })
|
||||
expect(mockGetWorkspaceBilledAccountUserId).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('does NOT fire for unlimited plans', () => {
|
||||
notifyTableRowUsage({ workspaceId: 'ws', currentRowCount: 0, addedRows: 10_000, limit: -1 })
|
||||
expect(mockGetWorkspaceBilledAccountUserId).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
*/
|
||||
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { maybeNotifyLimit } from '@/lib/billing/core/limit-notifications'
|
||||
import { getHighestPrioritySubscription } from '@/lib/billing/core/subscription'
|
||||
import { getPlanTypeForLimits } from '@/lib/billing/plan-helpers'
|
||||
import { getTablePlanLimits, type PlanName, type TablePlanLimits } from '@/lib/table/constants'
|
||||
@@ -12,6 +13,80 @@ import { getWorkspaceBilledAccountUserId } from '@/lib/workspaces/utils'
|
||||
|
||||
const logger = createLogger('TableBilling')
|
||||
|
||||
/** Warn band; an insert that crosses up into it (or up into 100%) triggers a notify. */
|
||||
const TABLE_WARN_PERCENT = 80
|
||||
const TABLE_REACHED_PERCENT = 100
|
||||
|
||||
/** Whether adding rows pushed the count UP across `threshold` (pre below, post at/above). */
|
||||
function crossedUp(prePct: number, postPct: number, threshold: number): boolean {
|
||||
return prePct < threshold && postPct >= threshold
|
||||
}
|
||||
|
||||
/**
|
||||
* Best-effort table row-limit email after an accepted insert. Resolves the
|
||||
* workspace's billed account, then delegates scope resolution + dedup + send to
|
||||
* {@link maybeNotifyLimit}. Never throws.
|
||||
*/
|
||||
async function maybeNotifyTableRowLimit(
|
||||
workspaceId: string,
|
||||
projectedRowCount: number,
|
||||
limit: number
|
||||
): Promise<void> {
|
||||
try {
|
||||
const billedUserId = await getWorkspaceBilledAccountUserId(workspaceId)
|
||||
if (!billedUserId) return
|
||||
|
||||
await maybeNotifyLimit({
|
||||
category: 'tables',
|
||||
billedUserId,
|
||||
workspaceId,
|
||||
currentUsage: projectedRowCount,
|
||||
limit,
|
||||
usageLabel: `${projectedRowCount.toLocaleString('en-US')} rows`,
|
||||
limitLabel: `${limit.toLocaleString('en-US')} rows`,
|
||||
})
|
||||
} catch (error) {
|
||||
logger.error('Error evaluating table row-limit notification:', error)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fire-and-forget the table row-limit threshold email for an accepted insert.
|
||||
* Edge-triggered: it only does any billing work when the write CROSSES UP into
|
||||
* the warn (80%) or reached (100%) band — not on every near-limit insert — so a
|
||||
* table sitting at e.g. 90% being inserted into pays nothing until it crosses
|
||||
* 100%. Shared by every insert path ({@link assertRowCapacity} and the
|
||||
* transactional upsert/import branches that check capacity with
|
||||
* {@link wouldExceedRowLimit} instead). Pass the pre-insert `currentRowCount` and
|
||||
* `addedRows` so both the pre and post percentages are known.
|
||||
*
|
||||
* Tables warn once per threshold and do not re-arm: a table that hit a threshold
|
||||
* then dropped (via deletes, which have no notify hook) won't re-warn on a later
|
||||
* climb. Deliberate trade-off — re-arm is storage-only (its single decrement
|
||||
* hook) to avoid per-delete billing-table reads. The atomic claim still dedups
|
||||
* concurrent crossings (no race).
|
||||
*/
|
||||
export function notifyTableRowUsage(params: {
|
||||
workspaceId: string
|
||||
currentRowCount: number
|
||||
addedRows: number
|
||||
limit: number
|
||||
}): void {
|
||||
if (params.limit <= 0) return
|
||||
const prePct = (params.currentRowCount / params.limit) * 100
|
||||
const postPct = ((params.currentRowCount + params.addedRows) / params.limit) * 100
|
||||
if (
|
||||
crossedUp(prePct, postPct, TABLE_WARN_PERCENT) ||
|
||||
crossedUp(prePct, postPct, TABLE_REACHED_PERCENT)
|
||||
) {
|
||||
void maybeNotifyTableRowLimit(
|
||||
params.workspaceId,
|
||||
params.currentRowCount + params.addedRows,
|
||||
params.limit
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Plan lookups hit billing + subscription tables (2-3 queries). Row-limit checks
|
||||
* run on every insert, so a short TTL keeps the hot path off the DB. Plan changes
|
||||
@@ -132,17 +207,23 @@ export function wouldExceedRowLimit(
|
||||
* connection (and locks) risks pool starvation. Callers already inside a tx should
|
||||
* fetch the limit up front and use {@link wouldExceedRowLimit} instead.
|
||||
*
|
||||
* Pure check (no side effects): returns the resolved limit so callers can fire
|
||||
* {@link notifyTableRowUsage} AFTER their insert commits — a pre-commit notify
|
||||
* would email (and burn the dedup claim) for a write that later rolls back.
|
||||
*
|
||||
* @returns the resolved plan row limit (-1 for unlimited)
|
||||
* @throws {TableRowLimitError} if `currentRowCount + addedRows` exceeds the limit
|
||||
*/
|
||||
export async function assertRowCapacity(params: {
|
||||
workspaceId: string
|
||||
currentRowCount: number
|
||||
addedRows: number
|
||||
}): Promise<void> {
|
||||
}): Promise<number> {
|
||||
const limit = await getMaxRowsPerTable(params.workspaceId)
|
||||
if (wouldExceedRowLimit(limit, params.currentRowCount, params.addedRows)) {
|
||||
throw new TableRowLimitError(limit)
|
||||
}
|
||||
return limit
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -9,7 +9,7 @@ import { userTableDefinitions, userTableRows } from '@sim/db/schema'
|
||||
import { createLogger } from '@sim/logger'
|
||||
import { generateId } from '@sim/utils/id'
|
||||
import { eq } from 'drizzle-orm'
|
||||
import { assertRowCapacity } from '@/lib/table/billing'
|
||||
import { assertRowCapacity, notifyTableRowUsage } from '@/lib/table/billing'
|
||||
import { CSV_MAX_BATCH_SIZE } from '@/lib/table/import'
|
||||
import { nKeysBetween } from '@/lib/table/order-key'
|
||||
import { acquireRowOrderLock } from '@/lib/table/rows/ordering'
|
||||
@@ -150,12 +150,12 @@ export async function importAppendRows(
|
||||
ctx: { workspaceId: string; userId?: string; requestId: string }
|
||||
): Promise<{ inserted: TableRow[]; table: TableDefinition }> {
|
||||
// Gate capacity before opening the tx — the lookup is a separate pool read.
|
||||
await assertRowCapacity({
|
||||
const rowLimit = await assertRowCapacity({
|
||||
workspaceId: ctx.workspaceId,
|
||||
currentRowCount: table.rowCount,
|
||||
addedRows: rows.length,
|
||||
})
|
||||
return db.transaction(async (trx) => {
|
||||
const result = await db.transaction(async (trx) => {
|
||||
let working = table
|
||||
if (additions.length > 0) {
|
||||
// Take the row-order lock before creating columns so this path uses the
|
||||
@@ -179,6 +179,13 @@ export async function importAppendRows(
|
||||
}
|
||||
return { inserted, table: working }
|
||||
})
|
||||
notifyTableRowUsage({
|
||||
workspaceId: ctx.workspaceId,
|
||||
currentRowCount: table.rowCount,
|
||||
addedRows: result.inserted.length,
|
||||
limit: rowLimit,
|
||||
})
|
||||
return result
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -193,12 +200,12 @@ export async function importReplaceRows(
|
||||
): Promise<ReplaceRowsResult> {
|
||||
// Replace deletes all existing rows, so the footprint is just the new set. Gate
|
||||
// before opening the tx — the plan lookup is a separate pool read.
|
||||
await assertRowCapacity({
|
||||
const rowLimit = await assertRowCapacity({
|
||||
workspaceId: data.workspaceId,
|
||||
currentRowCount: 0,
|
||||
addedRows: data.rows.length,
|
||||
})
|
||||
return db.transaction(async (trx) => {
|
||||
const result = await db.transaction(async (trx) => {
|
||||
let working = table
|
||||
if (additions.length > 0) {
|
||||
await acquireRowOrderLock(trx, table.id)
|
||||
@@ -211,4 +218,11 @@ export async function importReplaceRows(
|
||||
requestId
|
||||
)
|
||||
})
|
||||
notifyTableRowUsage({
|
||||
workspaceId: data.workspaceId,
|
||||
currentRowCount: 0,
|
||||
addedRows: result.insertedCount,
|
||||
limit: rowLimit,
|
||||
})
|
||||
return result
|
||||
}
|
||||
|
||||
@@ -17,7 +17,7 @@ import {
|
||||
type TableSchema,
|
||||
validateMapping,
|
||||
} from '@/lib/table'
|
||||
import { assertRowCapacity } from '@/lib/table/billing'
|
||||
import { assertRowCapacity, notifyTableRowUsage } from '@/lib/table/billing'
|
||||
import { withGeneratedColumnIds } from '@/lib/table/column-keys'
|
||||
import { appendTableEvent } from '@/lib/table/events'
|
||||
import {
|
||||
@@ -200,7 +200,7 @@ export async function runTableImport(payload: TableImportPayload): Promise<void>
|
||||
const owns = await updateJobProgress(tableId, inserted, importId)
|
||||
if (!owns) throw new ImportSupersededError()
|
||||
const coerced = coerceRowsForTable(rows, schema, headerToColumn)
|
||||
await assertRowCapacity({
|
||||
const rowLimit = await assertRowCapacity({
|
||||
workspaceId,
|
||||
currentRowCount: existingRowCount + inserted,
|
||||
addedRows: coerced.length,
|
||||
@@ -217,6 +217,12 @@ export async function runTableImport(payload: TableImportPayload): Promise<void>
|
||||
{ ...table, schema },
|
||||
requestId
|
||||
)
|
||||
notifyTableRowUsage({
|
||||
workspaceId,
|
||||
currentRowCount: existingRowCount + inserted,
|
||||
addedRows: result.inserted,
|
||||
limit: rowLimit,
|
||||
})
|
||||
inserted += result.inserted
|
||||
lastOrderKey = result.lastOrderKey
|
||||
// Emit after the first batch, then every interval, so the bar appears early without flooding.
|
||||
|
||||
@@ -20,6 +20,7 @@ import { isFeatureEnabled } from '@/lib/core/config/feature-flags'
|
||||
import {
|
||||
assertRowCapacity,
|
||||
getMaxRowsPerTable,
|
||||
notifyTableRowUsage,
|
||||
TableRowLimitError,
|
||||
wouldExceedRowLimit,
|
||||
} from '@/lib/table/billing'
|
||||
@@ -122,7 +123,7 @@ export async function insertRow(
|
||||
}
|
||||
|
||||
// Best-effort capacity check against the workspace's current plan limit.
|
||||
await assertRowCapacity({
|
||||
const rowLimit = await assertRowCapacity({
|
||||
workspaceId: table.workspaceId,
|
||||
currentRowCount: table.rowCount,
|
||||
addedRows: 1,
|
||||
@@ -143,6 +144,13 @@ export async function insertRow(
|
||||
now,
|
||||
})
|
||||
|
||||
notifyTableRowUsage({
|
||||
workspaceId: table.workspaceId,
|
||||
currentRowCount: table.rowCount,
|
||||
addedRows: 1,
|
||||
limit: rowLimit,
|
||||
})
|
||||
|
||||
logger.info(`[${requestId}] Inserted row ${rowId} into table ${data.tableId}`)
|
||||
|
||||
const insertedRow: TableRow = {
|
||||
@@ -193,13 +201,19 @@ export async function batchInsertRows(
|
||||
): Promise<TableRow[]> {
|
||||
// Best-effort capacity check against the workspace's current plan limit. Import
|
||||
// paths call `batchInsertRowsWithTx` directly and gate capacity up front instead.
|
||||
await assertRowCapacity({
|
||||
const rowLimit = await assertRowCapacity({
|
||||
workspaceId: table.workspaceId,
|
||||
currentRowCount: table.rowCount,
|
||||
addedRows: data.rows.length,
|
||||
})
|
||||
|
||||
const result = await db.transaction((trx) => batchInsertRowsWithTx(trx, data, table, requestId))
|
||||
notifyTableRowUsage({
|
||||
workspaceId: table.workspaceId,
|
||||
currentRowCount: table.rowCount,
|
||||
addedRows: result.length,
|
||||
limit: rowLimit,
|
||||
})
|
||||
dispatchAfterBatchInsert(table, result, requestId, data.userId)
|
||||
return result
|
||||
}
|
||||
@@ -348,12 +362,19 @@ export async function replaceTableRows(
|
||||
): Promise<ReplaceRowsResult> {
|
||||
// All existing rows are deleted, so the footprint is just the new set. Checked
|
||||
// before the tx opens — never inside it (the plan lookup is a separate pool read).
|
||||
await assertRowCapacity({
|
||||
const rowLimit = await assertRowCapacity({
|
||||
workspaceId: table.workspaceId,
|
||||
currentRowCount: 0,
|
||||
addedRows: data.rows.length,
|
||||
})
|
||||
return db.transaction((trx) => replaceTableRowsWithTx(trx, data, table, requestId))
|
||||
const result = await db.transaction((trx) => replaceTableRowsWithTx(trx, data, table, requestId))
|
||||
notifyTableRowUsage({
|
||||
workspaceId: table.workspaceId,
|
||||
currentRowCount: 0,
|
||||
addedRows: result.insertedCount,
|
||||
limit: rowLimit,
|
||||
})
|
||||
return result
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -673,6 +694,12 @@ export async function upsertRow(
|
||||
)
|
||||
|
||||
if (result.operation === 'insert') {
|
||||
notifyTableRowUsage({
|
||||
workspaceId: data.workspaceId,
|
||||
currentRowCount: table.rowCount,
|
||||
addedRows: 1,
|
||||
limit: rowLimit,
|
||||
})
|
||||
void fireTableTrigger(
|
||||
data.tableId,
|
||||
table.name,
|
||||
|
||||
@@ -15,7 +15,7 @@ import { generateId } from '@sim/utils/id'
|
||||
import { and, count, eq, isNull, sql } from 'drizzle-orm'
|
||||
import { generateRestoreName } from '@/lib/core/utils/restore-name'
|
||||
import type { DbOrTx } from '@/lib/db/types'
|
||||
import { assertRowCapacity } from '@/lib/table/billing'
|
||||
import { assertRowCapacity, notifyTableRowUsage } from '@/lib/table/billing'
|
||||
import { generateColumnId, getColumnId, withGeneratedColumnIds } from '@/lib/table/column-keys'
|
||||
import { COLUMN_TYPES, NAME_PATTERN, TABLE_LIMITS } from '@/lib/table/constants'
|
||||
import { EMPTY_JOB_FIELDS, latestJobForTable, latestJobsForTables } from '@/lib/table/jobs/service'
|
||||
@@ -285,8 +285,9 @@ export async function createTable(
|
||||
// Starter rows count against the plan too. Checked before the tx (the lookup is a
|
||||
// separate pool read) — a new table starts empty, so the footprint is just these.
|
||||
const initialRowCount = data.initialRowCount ?? 0
|
||||
let rowLimit: number | undefined
|
||||
if (initialRowCount > 0) {
|
||||
await assertRowCapacity({
|
||||
rowLimit = await assertRowCapacity({
|
||||
workspaceId: data.workspaceId,
|
||||
currentRowCount: 0,
|
||||
addedRows: initialRowCount,
|
||||
@@ -369,6 +370,15 @@ export async function createTable(
|
||||
throw error
|
||||
}
|
||||
|
||||
if (initialRowCount > 0 && rowLimit !== undefined) {
|
||||
notifyTableRowUsage({
|
||||
workspaceId: data.workspaceId,
|
||||
currentRowCount: 0,
|
||||
addedRows: initialRowCount,
|
||||
limit: rowLimit,
|
||||
})
|
||||
}
|
||||
|
||||
logger.info(`[${requestId}] Created table ${tableId} in workspace ${data.workspaceId}`)
|
||||
|
||||
return {
|
||||
|
||||
@@ -267,7 +267,7 @@ export async function uploadWorkspaceFile(
|
||||
)
|
||||
|
||||
try {
|
||||
await incrementStorageUsage(userId, fileBuffer.length)
|
||||
await incrementStorageUsage(userId, fileBuffer.length, workspaceId)
|
||||
} catch (storageError) {
|
||||
logger.error(`Failed to update storage tracking:`, storageError)
|
||||
}
|
||||
@@ -431,7 +431,7 @@ export async function registerUploadedWorkspaceFile(params: {
|
||||
}
|
||||
|
||||
try {
|
||||
await incrementStorageUsage(userId, verifiedSize)
|
||||
await incrementStorageUsage(userId, verifiedSize, workspaceId)
|
||||
} catch (storageError) {
|
||||
logger.error('Failed to update storage tracking:', storageError)
|
||||
}
|
||||
@@ -935,9 +935,9 @@ export async function updateWorkspaceFileContent(
|
||||
if (sizeDiff !== 0) {
|
||||
try {
|
||||
if (sizeDiff > 0) {
|
||||
await incrementStorageUsage(userId, sizeDiff)
|
||||
await incrementStorageUsage(userId, sizeDiff, workspaceId)
|
||||
} else {
|
||||
await decrementStorageUsage(userId, Math.abs(sizeDiff))
|
||||
await decrementStorageUsage(userId, Math.abs(sizeDiff), workspaceId)
|
||||
}
|
||||
} catch (storageError) {
|
||||
logger.error(`Failed to update storage tracking:`, storageError)
|
||||
|
||||
Reference in New Issue
Block a user