feat: add user level blob quota (#4114)
This commit is contained in:
10
Cargo.lock
generated
10
Cargo.lock
generated
@@ -1695,7 +1695,7 @@ dependencies = [
|
|||||||
[[package]]
|
[[package]]
|
||||||
name = "jwst"
|
name = "jwst"
|
||||||
version = "0.1.1"
|
version = "0.1.1"
|
||||||
source = "git+https://github.com/toeverything/OctoBase.git#b026b44f67e1043e83f31a21360a0baee122e819"
|
source = "git+https://github.com/toeverything/OctoBase.git#58f3bbdf97f391a535e772d32828a484376c4159"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"base64 0.21.3",
|
"base64 0.21.3",
|
||||||
@@ -1723,7 +1723,7 @@ dependencies = [
|
|||||||
[[package]]
|
[[package]]
|
||||||
name = "jwst-codec"
|
name = "jwst-codec"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
source = "git+https://github.com/toeverything/OctoBase.git#b026b44f67e1043e83f31a21360a0baee122e819"
|
source = "git+https://github.com/toeverything/OctoBase.git#58f3bbdf97f391a535e772d32828a484376c4159"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"arbitrary",
|
"arbitrary",
|
||||||
"bitvec",
|
"bitvec",
|
||||||
@@ -1743,7 +1743,7 @@ dependencies = [
|
|||||||
[[package]]
|
[[package]]
|
||||||
name = "jwst-logger"
|
name = "jwst-logger"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
source = "git+https://github.com/toeverything/OctoBase.git#b026b44f67e1043e83f31a21360a0baee122e819"
|
source = "git+https://github.com/toeverything/OctoBase.git#58f3bbdf97f391a535e772d32828a484376c4159"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"chrono",
|
"chrono",
|
||||||
"nu-ansi-term",
|
"nu-ansi-term",
|
||||||
@@ -1756,7 +1756,7 @@ dependencies = [
|
|||||||
[[package]]
|
[[package]]
|
||||||
name = "jwst-storage"
|
name = "jwst-storage"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
source = "git+https://github.com/toeverything/OctoBase.git#b026b44f67e1043e83f31a21360a0baee122e819"
|
source = "git+https://github.com/toeverything/OctoBase.git#58f3bbdf97f391a535e772d32828a484376c4159"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
@@ -1786,7 +1786,7 @@ dependencies = [
|
|||||||
[[package]]
|
[[package]]
|
||||||
name = "jwst-storage-migration"
|
name = "jwst-storage-migration"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
source = "git+https://github.com/toeverything/OctoBase.git#b026b44f67e1043e83f31a21360a0baee122e819"
|
source = "git+https://github.com/toeverything/OctoBase.git#58f3bbdf97f391a535e772d32828a484376c4159"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"sea-orm-migration",
|
"sea-orm-migration",
|
||||||
"tokio",
|
"tokio",
|
||||||
|
|||||||
@@ -186,6 +186,11 @@ export interface AFFiNEConfig {
|
|||||||
fs: {
|
fs: {
|
||||||
path: string;
|
path: string;
|
||||||
};
|
};
|
||||||
|
/**
|
||||||
|
* Free user storage quota
|
||||||
|
* @default 10 * 1024 * 1024 (10GB)
|
||||||
|
*/
|
||||||
|
quota: number;
|
||||||
};
|
};
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -55,6 +55,7 @@ export const getDefaultAFFiNEConfig: () => AFFiNEConfig = () => {
|
|||||||
AFFINE_SERVER_HOST: 'host',
|
AFFINE_SERVER_HOST: 'host',
|
||||||
AFFINE_SERVER_SUB_PATH: 'path',
|
AFFINE_SERVER_SUB_PATH: 'path',
|
||||||
AFFINE_ENV: 'affineEnv',
|
AFFINE_ENV: 'affineEnv',
|
||||||
|
AFFINE_FREE_USER_QUOTA: 'objectStorage.quota',
|
||||||
DATABASE_URL: 'db.url',
|
DATABASE_URL: 'db.url',
|
||||||
ENABLE_R2_OBJECT_STORAGE: ['objectStorage.r2.enabled', 'boolean'],
|
ENABLE_R2_OBJECT_STORAGE: ['objectStorage.r2.enabled', 'boolean'],
|
||||||
R2_OBJECT_STORAGE_ACCOUNT_ID: 'objectStorage.r2.accountId',
|
R2_OBJECT_STORAGE_ACCOUNT_ID: 'objectStorage.r2.accountId',
|
||||||
@@ -170,6 +171,7 @@ export const getDefaultAFFiNEConfig: () => AFFiNEConfig = () => {
|
|||||||
fs: {
|
fs: {
|
||||||
path: join(homedir(), '.affine-storage'),
|
path: join(homedir(), '.affine-storage'),
|
||||||
},
|
},
|
||||||
|
quota: 10 * 1024 * 1024,
|
||||||
},
|
},
|
||||||
rateLimiter: {
|
rateLimiter: {
|
||||||
ttl: 60,
|
ttl: 60,
|
||||||
|
|||||||
@@ -27,6 +27,7 @@ import type { User, Workspace } from '@prisma/client';
|
|||||||
import GraphQLUpload from 'graphql-upload/GraphQLUpload.mjs';
|
import GraphQLUpload from 'graphql-upload/GraphQLUpload.mjs';
|
||||||
import { applyUpdate, Doc } from 'yjs';
|
import { applyUpdate, Doc } from 'yjs';
|
||||||
|
|
||||||
|
import { Config } from '../../config';
|
||||||
import { PrismaService } from '../../prisma';
|
import { PrismaService } from '../../prisma';
|
||||||
import { StorageProvide } from '../../storage';
|
import { StorageProvide } from '../../storage';
|
||||||
import { CloudThrottlerGuard, Throttle } from '../../throttler';
|
import { CloudThrottlerGuard, Throttle } from '../../throttler';
|
||||||
@@ -130,6 +131,7 @@ export class UpdateWorkspaceInput extends PickType(
|
|||||||
export class WorkspaceResolver {
|
export class WorkspaceResolver {
|
||||||
constructor(
|
constructor(
|
||||||
private readonly auth: AuthService,
|
private readonly auth: AuthService,
|
||||||
|
private readonly config: Config,
|
||||||
private readonly mailer: MailService,
|
private readonly mailer: MailService,
|
||||||
private readonly prisma: PrismaService,
|
private readonly prisma: PrismaService,
|
||||||
private readonly permissionProvider: PermissionService,
|
private readonly permissionProvider: PermissionService,
|
||||||
@@ -610,17 +612,28 @@ export class WorkspaceResolver {
|
|||||||
) {
|
) {
|
||||||
await this.permissionProvider.check(workspaceId, user.id);
|
await this.permissionProvider.check(workspaceId, user.id);
|
||||||
|
|
||||||
return this.storage.blobsSize(workspaceId).then(size => ({ size }));
|
return this.storage.blobsSize([workspaceId]).then(size => ({ size }));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Query(() => WorkspaceBlobSizes)
|
@Query(() => WorkspaceBlobSizes)
|
||||||
async collectAllBlobSizes(@CurrentUser() user: User) {
|
async collectAllBlobSizes(@CurrentUser() user: UserType) {
|
||||||
const workspaces = await this.workspaces(user);
|
const workspaces = await this.prisma.userWorkspacePermission
|
||||||
|
.findMany({
|
||||||
const size = (
|
where: {
|
||||||
await Promise.all(workspaces.map(({ id }) => this.storage.blobsSize(id)))
|
userId: user.id,
|
||||||
).reduce((prev, curr) => prev + curr, 0);
|
accepted: true,
|
||||||
|
},
|
||||||
|
select: {
|
||||||
|
workspace: {
|
||||||
|
select: {
|
||||||
|
id: true,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
.then(data => data.map(({ workspace }) => workspace.id));
|
||||||
|
|
||||||
|
const size = await this.storage.blobsSize(workspaces);
|
||||||
return { size };
|
return { size };
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -632,6 +645,12 @@ export class WorkspaceResolver {
|
|||||||
blob: FileUpload
|
blob: FileUpload
|
||||||
) {
|
) {
|
||||||
await this.permissionProvider.check(workspaceId, user.id, Permission.Write);
|
await this.permissionProvider.check(workspaceId, user.id, Permission.Write);
|
||||||
|
const quota = this.config.objectStorage.quota;
|
||||||
|
const { size } = await this.collectAllBlobSizes(user);
|
||||||
|
|
||||||
|
if (size > quota) {
|
||||||
|
throw new ForbiddenException('storage size limit exceeded');
|
||||||
|
}
|
||||||
|
|
||||||
const buffer = await new Promise<Buffer>((resolve, reject) => {
|
const buffer = await new Promise<Buffer>((resolve, reject) => {
|
||||||
const stream = blob.createReadStream();
|
const stream = blob.createReadStream();
|
||||||
@@ -645,6 +664,10 @@ export class WorkspaceResolver {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
if (size + buffer.length > quota) {
|
||||||
|
throw new ForbiddenException('storage size limit exceeded');
|
||||||
|
}
|
||||||
|
|
||||||
return this.storage.uploadBlob(workspaceId, buffer);
|
return this.storage.uploadBlob(workspaceId, buffer);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -339,7 +339,6 @@ async function collectBlobSizes(
|
|||||||
const res = await request(app.getHttpServer())
|
const res = await request(app.getHttpServer())
|
||||||
.post(gql)
|
.post(gql)
|
||||||
.auth(token, { type: 'bearer' })
|
.auth(token, { type: 'bearer' })
|
||||||
.set({ 'x-request-id': 'test', 'x-operation-name': 'test' })
|
|
||||||
.send({
|
.send({
|
||||||
query: `
|
query: `
|
||||||
query {
|
query {
|
||||||
@@ -353,6 +352,26 @@ async function collectBlobSizes(
|
|||||||
return res.body.data.collectBlobSizes.size;
|
return res.body.data.collectBlobSizes.size;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async function collectAllBlobSizes(
|
||||||
|
app: INestApplication,
|
||||||
|
token: string
|
||||||
|
): Promise<number> {
|
||||||
|
const res = await request(app.getHttpServer())
|
||||||
|
.post(gql)
|
||||||
|
.auth(token, { type: 'bearer' })
|
||||||
|
.send({
|
||||||
|
query: `
|
||||||
|
query {
|
||||||
|
collectAllBlobSizes {
|
||||||
|
size
|
||||||
|
}
|
||||||
|
}
|
||||||
|
`,
|
||||||
|
})
|
||||||
|
.expect(200);
|
||||||
|
return res.body.data.collectAllBlobSizes.size;
|
||||||
|
}
|
||||||
|
|
||||||
async function setBlob(
|
async function setBlob(
|
||||||
app: INestApplication,
|
app: INestApplication,
|
||||||
token: string,
|
token: string,
|
||||||
@@ -447,6 +466,7 @@ async function getInviteInfo(
|
|||||||
export {
|
export {
|
||||||
acceptInvite,
|
acceptInvite,
|
||||||
acceptInviteById,
|
acceptInviteById,
|
||||||
|
collectAllBlobSizes,
|
||||||
collectBlobSizes,
|
collectBlobSizes,
|
||||||
createTestApp,
|
createTestApp,
|
||||||
createWorkspace,
|
createWorkspace,
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import request from 'supertest';
|
|||||||
|
|
||||||
import { AppModule } from '../app';
|
import { AppModule } from '../app';
|
||||||
import {
|
import {
|
||||||
|
collectAllBlobSizes,
|
||||||
collectBlobSizes,
|
collectBlobSizes,
|
||||||
createWorkspace,
|
createWorkspace,
|
||||||
listBlobs,
|
listBlobs,
|
||||||
@@ -108,4 +109,25 @@ describe('Workspace Module - Blobs', () => {
|
|||||||
const size = await collectBlobSizes(app, u1.token.token, workspace.id);
|
const size = await collectBlobSizes(app, u1.token.token, workspace.id);
|
||||||
ok(size === 4, 'failed to collect blob sizes');
|
ok(size === 4, 'failed to collect blob sizes');
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('should calc all blobs size', async () => {
|
||||||
|
const u1 = await signUp(app, 'u1', 'u1@affine.pro', '1');
|
||||||
|
|
||||||
|
const workspace1 = await createWorkspace(app, u1.token.token);
|
||||||
|
|
||||||
|
const buffer1 = Buffer.from([0, 0]);
|
||||||
|
await setBlob(app, u1.token.token, workspace1.id, buffer1);
|
||||||
|
const buffer2 = Buffer.from([0, 1]);
|
||||||
|
await setBlob(app, u1.token.token, workspace1.id, buffer2);
|
||||||
|
|
||||||
|
const workspace2 = await createWorkspace(app, u1.token.token);
|
||||||
|
|
||||||
|
const buffer3 = Buffer.from([0, 0]);
|
||||||
|
await setBlob(app, u1.token.token, workspace2.id, buffer3);
|
||||||
|
const buffer4 = Buffer.from([0, 1]);
|
||||||
|
await setBlob(app, u1.token.token, workspace2.id, buffer4);
|
||||||
|
|
||||||
|
const size = await collectAllBlobSizes(app, u1.token.token);
|
||||||
|
ok(size === 8, 'failed to collect all blob sizes');
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
7
packages/storage/index.d.ts
vendored
7
packages/storage/index.d.ts
vendored
@@ -16,7 +16,7 @@ export class Storage {
|
|||||||
/** Delete a blob from workspace storage. */
|
/** Delete a blob from workspace storage. */
|
||||||
deleteBlob(workspaceId: string, hash: string): Promise<boolean>;
|
deleteBlob(workspaceId: string, hash: string): Promise<boolean>;
|
||||||
/** Workspace size taken by blobs. */
|
/** Workspace size taken by blobs. */
|
||||||
blobsSize(workspaceId: string): Promise<number>;
|
blobsSize(workspaces: Array<string>): Promise<number>;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface Blob {
|
export interface Blob {
|
||||||
@@ -26,5 +26,8 @@ export interface Blob {
|
|||||||
data: Buffer;
|
data: Buffer;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Merge updates in form like `Y.applyUpdate(doc, update)` way and return the result binary. */
|
/**
|
||||||
|
* Merge updates in form like `Y.applyUpdate(doc, update)` way and return the
|
||||||
|
* result binary.
|
||||||
|
*/
|
||||||
export function mergeUpdatesInApplyWay(updates: Array<Buffer>): Buffer;
|
export function mergeUpdatesInApplyWay(updates: Array<Buffer>): Buffer;
|
||||||
|
|||||||
@@ -9,7 +9,6 @@ use std::{
|
|||||||
use jwst::BlobStorage;
|
use jwst::BlobStorage;
|
||||||
use jwst_codec::Doc;
|
use jwst_codec::Doc;
|
||||||
use jwst_storage::{BlobStorageType, JwstStorage, JwstStorageError};
|
use jwst_storage::{BlobStorageType, JwstStorage, JwstStorageError};
|
||||||
|
|
||||||
use napi::{bindgen_prelude::*, Error, Result, Status};
|
use napi::{bindgen_prelude::*, Error, Result, Status};
|
||||||
|
|
||||||
#[macro_use]
|
#[macro_use]
|
||||||
@@ -112,7 +111,11 @@ impl Storage {
|
|||||||
(id, ext.map(|ext| HashMap::from([("format".into(), ext)])))
|
(id, ext.map(|ext| HashMap::from([("format".into(), ext)])))
|
||||||
};
|
};
|
||||||
|
|
||||||
let Ok(meta) = self.blobs().get_metadata(Some(workspace_id.clone()), id.clone(), params.clone()).await else {
|
let Ok(meta) = self
|
||||||
|
.blobs()
|
||||||
|
.get_metadata(Some(workspace_id.clone()), id.clone(), params.clone())
|
||||||
|
.await
|
||||||
|
else {
|
||||||
return Ok(None);
|
return Ok(None);
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -144,12 +147,13 @@ impl Storage {
|
|||||||
|
|
||||||
/// Workspace size taken by blobs.
|
/// Workspace size taken by blobs.
|
||||||
#[napi]
|
#[napi]
|
||||||
pub async fn blobs_size(&self, workspace_id: String) -> Result<i64> {
|
pub async fn blobs_size(&self, workspaces: Vec<String>) -> Result<i64> {
|
||||||
map_err!(self.blobs().get_blobs_size(workspace_id).await)
|
map_err!(self.blobs().get_blobs_size(workspaces).await)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Merge updates in form like `Y.applyUpdate(doc, update)` way and return the result binary.
|
/// Merge updates in form like `Y.applyUpdate(doc, update)` way and return the
|
||||||
|
/// result binary.
|
||||||
#[napi(catch_unwind)]
|
#[napi(catch_unwind)]
|
||||||
pub fn merge_updates_in_apply_way(updates: Vec<Buffer>) -> Result<Buffer> {
|
pub fn merge_updates_in_apply_way(updates: Vec<Buffer>) -> Result<Buffer> {
|
||||||
let mut doc = Doc::default();
|
let mut doc = Doc::default();
|
||||||
|
|||||||
Reference in New Issue
Block a user