mirror of
https://github.com/yusufipk/OpenFrame.git
synced 2026-09-11 17:46:06 +00:00
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.
541 lines
17 KiB
TypeScript
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 }
|
|
);
|