Files
OpenFrame/lib/admin-stats.ts
T
yusufipek 4ff801738c fix(uploads): count a Bunny upload from the moment it is admitted
A Bunny init asked the quota whether it could store zero bytes, which is a
question with only one answer. Nothing an upload was about to consume was
visible to the next request, so every init inside the same window read the
same total and every one of them passed, and an upload that could never
fit was only refused after it had been sent.

The client now declares the size up front. It is checked against the
account's remaining room before Bunny is asked for anything, and held as
a reservation the next init has to see. The declaration is a claim rather
than proof, so it is signed into the upload token: the same token already
binds the video id, which is what makes the reservation safe to release
on a caller's say-so, since releasing it costs them the video it belongs
to.

The declared size is then written onto the version or asset row and the
reservation is dropped in the same transaction, because Bunny reports no
size at all for a video until it has finished encoding it. On a half hour
of footage that is most of an hour during which the upload did not appear
on the uploader's own storage page and did not count against the next
upload. Per-video accounting now takes the larger of what Bunny reports
and what was declared, so the estimate stands in until the real figure
arrives and Bunny's wins once it does.

Two smaller things came out of the same reading. The asset route's
in-transaction fallback compared against the plan limit, so a caller
quoting a reservation that no longer existed was measured against 200 GiB
even on a trial worth three. And the guest branch reserves without being
able to release early, because a guest grant is bound to our video id and
the caller's network context rather than to the Bunny video, which would
let the reservation be dropped while the upload it stands for carried on.
2026-08-18 10:35:08 +03:00

541 lines
17 KiB
TypeScript

import { unstable_cache } from 'next/cache';
import { db } from '@/lib/db';
import { r2Client, R2_BUCKET_NAME } from '@/lib/r2';
import { ListObjectsV2Command, type ListObjectsV2CommandInput } from '@aws-sdk/client-s3';
import { isBunnyUploadsEnabled, isStripeBillingEnabled } from '@/lib/feature-flags';
import { getStripe, getStripePriceId } from '@/lib/stripe';
import { logError } from '@/lib/logger';
const BUNNY_API_BASE = 'https://video.bunnycdn.com';
const STORAGE_CACHE_SECONDS = 120;
interface R2StorageSnapshot {
fileSizes: Map<string, number>;
totalBytes: number;
refreshedAt: string;
}
const globalForAdminStats = globalThis as unknown as {
adminR2StorageSnapshot?: R2StorageSnapshot;
adminR2StorageSnapshotPromise?: Promise<R2StorageSnapshot>;
};
interface BunnyStorageStats {
totalBytes: number;
byVideoId: Record<string, number>;
}
function bigintToNumber(value: bigint): number {
return value > BigInt(Number.MAX_SAFE_INTEGER) ? Number.MAX_SAFE_INTEGER : Number(value);
}
function getBunnyConfig(): { apiKey: string; libraryId: string } {
const apiKey = process.env.BUNNY_STREAM_API_KEY;
const libraryId =
process.env.BUNNY_STREAM_LIBRARY_ID || process.env.NEXT_PUBLIC_BUNNY_STREAM_LIBRARY_ID;
if (!apiKey || !libraryId) {
throw new Error('Missing Bunny Stream credentials.');
}
return { apiKey, libraryId };
}
function toRecord(value: unknown): Record<string, unknown> | null {
if (!value || typeof value !== 'object') return null;
return value as Record<string, unknown>;
}
function parseBunnyVideoStorageBytes(item: unknown): number {
const record = toRecord(item);
if (!record) return 0;
const candidates = ['storageSize', 'storage', 'size'];
for (const key of candidates) {
const value = record[key];
if (typeof value === 'number' && Number.isFinite(value) && value > 0) {
return value;
}
}
return 0;
}
function parseBunnyVideoGuid(item: unknown): string | null {
const record = toRecord(item);
if (!record) return null;
const value = record.guid;
return typeof value === 'string' && value.length > 0 ? value : null;
}
async function listAllR2FileSizes(): Promise<Map<string, number>> {
const fileSizes = new Map<string, number>();
let isTruncated = true;
let continuationToken: string | undefined;
while (isTruncated) {
const commandParams: ListObjectsV2CommandInput = { Bucket: R2_BUCKET_NAME };
if (continuationToken) {
commandParams.ContinuationToken = continuationToken;
}
const data = await r2Client.send(new ListObjectsV2Command(commandParams));
if (data.Contents) {
for (const item of data.Contents) {
if (item.Key) fileSizes.set(item.Key, item.Size || 0);
}
}
isTruncated = data.IsTruncated ?? false;
continuationToken = data.NextContinuationToken;
}
return fileSizes;
}
async function buildR2StorageSnapshot(): Promise<R2StorageSnapshot> {
const fileSizes = await listAllR2FileSizes();
let totalBytes = 0;
for (const size of fileSizes.values()) {
totalBytes += size;
}
return {
fileSizes,
totalBytes,
refreshedAt: new Date().toISOString(),
};
}
async function getR2StorageSnapshot(): Promise<R2StorageSnapshot> {
if (globalForAdminStats.adminR2StorageSnapshot) {
return globalForAdminStats.adminR2StorageSnapshot;
}
return Promise.reject(
new Error(
'R2 storage snapshot is not available. Trigger a manual refresh from admin dashboard.'
)
);
}
export async function refreshR2StorageSnapshot(): Promise<string> {
// Single-flight. The promise slot was declared and cleared but never read, so two
// concurrent admin refreshes each walked the whole bucket. A second caller now joins
// the walk already in progress.
const inFlight = globalForAdminStats.adminR2StorageSnapshotPromise;
if (inFlight) {
return (await inFlight).refreshedAt;
}
const pending = buildR2StorageSnapshot();
globalForAdminStats.adminR2StorageSnapshotPromise = pending;
try {
const snapshot = await pending;
globalForAdminStats.adminR2StorageSnapshot = snapshot;
return snapshot.refreshedAt;
} finally {
globalForAdminStats.adminR2StorageSnapshotPromise = undefined;
}
}
async function fetchBunnyStorageStats(): Promise<BunnyStorageStats> {
// isBunnyUploadsEnabled(), not isBunnyUploadsFeatureEnabled(): the flag alone defaults
// to on, so a self-hosted install that never configured Bunny threw
// "Missing Bunny Stream credentials." out of getBunnyConfig() below and the dashboard
// reported -1 instead of zero.
if (!isBunnyUploadsEnabled()) {
return { totalBytes: 0, byVideoId: {} };
}
const { apiKey, libraryId } = getBunnyConfig();
const byVideoId: Record<string, number> = {};
let totalBytes = 0;
let page = 1;
const itemsPerPage = 100;
while (page <= 200) {
const response = await fetch(
`${BUNNY_API_BASE}/library/${libraryId}/videos?page=${page}&itemsPerPage=${itemsPerPage}`,
{ headers: { AccessKey: apiKey }, cache: 'no-store' }
);
if (!response.ok) {
throw new Error(`Bunny API failed (${response.status})`);
}
const json = await response.json();
const record = toRecord(json);
if (!record) break;
const rawItems = Array.isArray(record.items)
? record.items
: Array.isArray(record.Items)
? record.Items
: [];
if (rawItems.length === 0) break;
for (const rawItem of rawItems) {
const guid = parseBunnyVideoGuid(rawItem);
if (!guid) continue;
const storageBytes = parseBunnyVideoStorageBytes(rawItem);
byVideoId[guid] = storageBytes;
totalBytes += storageBytes;
}
const totalItems =
typeof record.totalItems === 'number'
? record.totalItems
: typeof record.TotalItems === 'number'
? record.TotalItems
: null;
if (totalItems !== null && page * itemsPerPage >= totalItems) {
break;
}
page += 1;
}
return { totalBytes, byVideoId };
}
export async function getCachedTotalStorage(): Promise<number> {
try {
const snapshot = await getR2StorageSnapshot();
return snapshot.totalBytes;
} catch (err) {
logError('Failed to fetch total storage stats:', err);
return -1;
}
}
export const getCachedBunnyStorageStats = unstable_cache(
async () => {
try {
return await fetchBunnyStorageStats();
} catch (err) {
logError('Failed to fetch Bunny storage stats:', err);
return { totalBytes: -1, byVideoId: {} } as BunnyStorageStats;
}
},
['admin-bunny-storage'],
{ revalidate: STORAGE_CACHE_SECONDS }
);
export const getCachedUserBunnyStorage = unstable_cache(
async () => {
const perUserStorage: Record<string, number> = {};
try {
const bunnyStats = await getCachedBunnyStorageStats();
if (bunnyStats.totalBytes < 0) return perUserStorage;
const [bunnyVersions, bunnyAssets] = await Promise.all([
db.videoVersion.findMany({
where: { providerId: 'bunny' },
select: {
videoId: true,
// What the uploader declared, used as a floor below.
sizeBytes: true,
video: {
select: {
project: {
// The workspace owner, not the project owner. lib/storage-quota.ts bills
// R2 versions to the workspace owner and comment media below does the
// same, and getCachedUserBunnyStorage feeds getUserTotalStorageBytes, so
// the moment project and workspace ownership can differ one workspace's
// Bunny bytes and its R2 bytes would count against two different quotas.
select: { workspace: { select: { ownerId: true } } },
},
},
},
},
}),
db.videoAsset.findMany({
where: {
provider: 'BUNNY',
providerVideoId: { not: null },
},
select: {
providerVideoId: true,
billedUserId: true,
sizeBytes: true,
},
}),
]);
/**
* What this video costs us, as the larger of the two numbers we have.
*
* Bunny reports nothing for a video until it has finished encoding it,
* which on a half-hour source is most of an hour, and reading that zero
* literally meant an upload was free for as long as it was being
* processed: it did not show on the uploader's storage page and it did not
* count against the next upload's quota check. The size declared when the
* upload was admitted stands in until Bunny has a figure of its own, and
* Bunny's wins once it arrives, because the renditions it makes are the
* real bill and they are larger than the source.
*/
const chargeableSize = (reported: number, declared: bigint | null): number => {
const declaredBytes = declared === null ? 0 : Number(declared);
return reported > declaredBytes ? reported : declaredBytes;
};
const seenVideoIds = new Set<string>();
for (const version of bunnyVersions) {
const ownerId = version.video.project.workspace.ownerId;
const dedupeKey = `${ownerId}:${version.videoId}`;
if (seenVideoIds.has(dedupeKey)) continue;
seenVideoIds.add(dedupeKey);
const size = chargeableSize(bunnyStats.byVideoId[version.videoId] || 0, version.sizeBytes);
perUserStorage[ownerId] = (perUserStorage[ownerId] || 0) + size;
}
for (const asset of bunnyAssets) {
if (!asset.providerVideoId) continue;
const billedUserId = asset.billedUserId;
const dedupeKey = `${billedUserId}:${asset.providerVideoId}`;
if (seenVideoIds.has(dedupeKey)) continue;
seenVideoIds.add(dedupeKey);
const size = chargeableSize(
bunnyStats.byVideoId[asset.providerVideoId] || 0,
asset.sizeBytes
);
perUserStorage[billedUserId] = (perUserStorage[billedUserId] || 0) + size;
}
} catch (err) {
logError('Failed to calculate per-user Bunny storage:', err);
}
return perUserStorage;
},
['admin-user-bunny-storage'],
{ revalidate: STORAGE_CACHE_SECONDS }
);
export async function getCachedUserMediaStorage(): Promise<
Record<string, { total: number; voice: number; image: number }>
> {
// Return a plain object so it maps cleanly out of server component boundaries
const userStorage: Record<string, { total: number; voice: number; image: number }> = {};
try {
const snapshot = await getR2StorageSnapshot();
const seenKeys = new Set<string>();
const [mediaComments, imageAssets, audioAssets] = await Promise.all([
db.comment.findMany({
where: { OR: [{ voiceUrl: { not: null } }, { imageUrl: { not: null } }] },
select: {
voiceUrl: true,
imageUrl: true,
version: {
select: {
video: {
select: {
project: {
select: {
workspace: {
select: { ownerId: true },
},
},
},
},
},
},
},
},
}),
db.videoAsset.findMany({
where: { provider: 'R2_IMAGE' },
select: {
sourceUrl: true,
billedUserId: true,
},
}),
db.videoAsset.findMany({
where: { provider: 'R2_AUDIO' },
select: {
sourceUrl: true,
billedUserId: true,
},
}),
]);
for (const comment of mediaComments) {
const billedUserId = comment.version.video.project.workspace.ownerId;
if (!billedUserId) continue;
if (!userStorage[billedUserId]) {
userStorage[billedUserId] = { total: 0, voice: 0, image: 0 };
}
if (comment.voiceUrl) {
const keyParts = comment.voiceUrl.split('/');
const filename = keyParts[keyParts.length - 1];
const r2Key = `voice/${filename}`;
const dedupeKey = `${billedUserId}:${r2Key}`;
if (!seenKeys.has(dedupeKey)) {
seenKeys.add(dedupeKey);
const size = snapshot.fileSizes.get(r2Key) || 0;
userStorage[billedUserId].voice += size;
userStorage[billedUserId].total += size;
}
}
if (comment.imageUrl) {
const keyParts = comment.imageUrl.split('/');
const filename = keyParts[keyParts.length - 1];
const r2Key = `images/${filename}`;
const dedupeKey = `${billedUserId}:${r2Key}`;
if (!seenKeys.has(dedupeKey)) {
seenKeys.add(dedupeKey);
const size = snapshot.fileSizes.get(r2Key) || 0;
userStorage[billedUserId].image += size;
userStorage[billedUserId].total += size;
}
}
}
for (const asset of imageAssets) {
const billedUserId = asset.billedUserId;
if (!billedUserId) continue;
if (!userStorage[billedUserId]) {
userStorage[billedUserId] = { total: 0, voice: 0, image: 0 };
}
const keyParts = asset.sourceUrl.split('/');
const filename = keyParts[keyParts.length - 1];
if (!filename) continue;
const r2Key = `images/${filename}`;
const dedupeKey = `${billedUserId}:${r2Key}`;
if (seenKeys.has(dedupeKey)) continue;
seenKeys.add(dedupeKey);
const size = snapshot.fileSizes.get(r2Key) || 0;
userStorage[billedUserId].image += size;
userStorage[billedUserId].total += size;
}
for (const asset of audioAssets) {
const billedUserId = asset.billedUserId;
if (!billedUserId) continue;
if (!userStorage[billedUserId]) {
userStorage[billedUserId] = { total: 0, voice: 0, image: 0 };
}
const keyParts = asset.sourceUrl.split('/');
const filename = keyParts[keyParts.length - 1];
if (!filename) continue;
const r2Key = `voice/${filename}`;
const dedupeKey = `${billedUserId}:${r2Key}`;
if (seenKeys.has(dedupeKey)) continue;
seenKeys.add(dedupeKey);
const size = snapshot.fileSizes.get(r2Key) || 0;
userStorage[billedUserId].voice += size;
userStorage[billedUserId].total += size;
}
} catch (err) {
logError('Failed to parse user storage:', err);
}
return userStorage;
}
export const getCachedUserDownloadEgress = unstable_cache(
async () => {
const perUserDownloadEgress: Record<string, number> = {};
try {
const grouped = await db.downloadEgressEvent.groupBy({
by: ['billedUserId'],
_sum: {
estimatedBytes: true,
},
});
for (const row of grouped) {
perUserDownloadEgress[row.billedUserId] = row._sum.estimatedBytes
? bigintToNumber(row._sum.estimatedBytes)
: 0;
}
} catch (err) {
logError('Failed to calculate per-user download egress:', err);
}
return perUserDownloadEgress;
},
['admin-user-download-egress'],
{ revalidate: STORAGE_CACHE_SECONDS }
);
export interface StripeStats {
activeSubscribers: number;
trialingUsers: number;
pastDueUsers: number;
canceledUsers: number;
freeUsers: number;
/** UNPAID, INCOMPLETE and INCOMPLETE_EXPIRED, which belong to none of the buckets above. */
otherStatusUsers: number;
mrrCents: number;
currency: string;
}
const STRIPE_STATS_CACHE_SECONDS = 300;
export const getCachedStripeStats = unstable_cache(
async (): Promise<StripeStats | null> => {
if (!isStripeBillingEnabled()) return null;
try {
const statusCounts = await db.user.groupBy({
by: ['subscriptionStatus'],
_count: { id: true },
});
const counts: Record<string, number> = {};
for (const row of statusCounts) {
counts[row.subscriptionStatus] = row._count.id;
}
const activeSubscribers = counts['ACTIVE'] ?? 0;
const trialingUsers = counts['TRIALING'] ?? 0;
const pastDueUsers = counts['PAST_DUE'] ?? 0;
const canceledUsers = counts['CANCELED'] ?? 0;
const freeUsers = counts['FREE'] ?? 0;
// UNPAID, INCOMPLETE and INCOMPLETE_EXPIRED belonged to none of the five buckets
// above, so those users were counted nowhere and the totals silently did not add
// up to the user table.
const otherStatusUsers =
(counts['UNPAID'] ?? 0) + (counts['INCOMPLETE'] ?? 0) + (counts['INCOMPLETE_EXPIRED'] ?? 0);
let mrrCents = 0;
let currency = 'usd';
try {
const stripe = getStripe();
const priceId = getStripePriceId();
const price = await stripe.prices.retrieve(priceId);
const unitAmount = price.unit_amount ?? 0;
currency = price.currency ?? 'usd';
mrrCents = activeSubscribers * unitAmount;
} catch (err) {
logError('Failed to fetch Stripe price for MRR calculation:', err);
}
return {
activeSubscribers,
trialingUsers,
pastDueUsers,
canceledUsers,
freeUsers,
otherStatusUsers,
mrrCents,
currency,
};
} catch (err) {
logError('Failed to fetch Stripe stats:', err);
return null;
}
},
['admin-stripe-stats'],
{ revalidate: STRIPE_STATS_CACHE_SECONDS }
);