Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 49 additions & 0 deletions lib/core/bucketEntries/MongoDBBucketEntriesRepository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,21 @@ import { BucketEntry, BucketEntryWithFrame } from './BucketEntry';
import { BucketEntriesRepository } from './Repository';
import { ObjectId } from 'mongodb';

/**
* The fields that mark an entry as backed by storage: `frame` on v1 entries,
* `index` and `hmac` on v2 uploads. Gateway entries (createEntry) set none of
* them, which is what makes them safe to drop wholesale.
*/
const SHARD_MARKERS = ['frame', 'index', 'hmac.value'];

const METADATA_ONLY = {
$and: SHARD_MARKERS.map((field) => ({ [field]: { $exists: false } })),
};

const SHARD_BACKED = {
$or: SHARD_MARKERS.map((field) => ({ [field]: { $exists: true } })),
};

interface BucketEntryModel extends Omit<BucketEntry, 'id'> {
_id: string;
created: Date;
Expand Down Expand Up @@ -90,6 +105,40 @@ export class MongoDBBucketEntriesRepository implements BucketEntriesRepository {
return bucketEntries.map(formatFromMongoToBucketEntry);
}

async hasEntriesByBucket(bucketId: string): Promise<boolean> {
const found = await this.model
.countDocuments({ bucket: bucketId }, { limit: 1 })
.read('primary')
Comment thread
sg-gs marked this conversation as resolved.
.exec();

return found > 0;
}

async hasShardBackedEntriesByBucket(bucketId: string): Promise<boolean> {
const found = await this.model
.countDocuments({ bucket: bucketId, ...SHARD_BACKED }, { limit: 1 })
.read('primary')
.exec();

return found > 0;
}

async sumMetadataOnlyBytesByBucket(bucketId: string): Promise<number> {
const [result] = await this.model
.aggregate([
Comment thread
sg-gs marked this conversation as resolved.
{ $match: { bucket: new ObjectId(bucketId), ...METADATA_ONLY } },
{ $group: { _id: null, bytes: { $sum: '$size' } } },
])
.read('primary')
.exec();

return result?.bytes ?? 0;
}

async deleteMetadataOnlyByBucket(bucketId: string): Promise<void> {
await this.model.deleteMany({ bucket: bucketId, ...METADATA_ONLY });
}

async findByIds(ids: string[]): Promise<BucketEntry[]> {
const bucketEntries = await this.model.find({ _id: { $in: ids } });

Expand Down
4 changes: 4 additions & 0 deletions lib/core/bucketEntries/Repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,10 @@ export interface BucketEntriesRepository {
count(where: Partial<BucketEntry>): Promise<number>;
findOne(where: Partial<BucketEntry>): Promise<BucketEntry | null>;
findByBucket(bucketId: Bucket['id'], limit: number, offset: number): Promise<BucketEntry[]>;
hasEntriesByBucket(bucketId: Bucket['id']): Promise<boolean>;
hasShardBackedEntriesByBucket(bucketId: Bucket['id']): Promise<boolean>;
sumMetadataOnlyBytesByBucket(bucketId: Bucket['id']): Promise<number>;
deleteMetadataOnlyByBucket(bucketId: Bucket['id']): Promise<void>;
findByIds(ids: BucketEntry['id'][]): Promise<BucketEntry[]>;
findOneWithFrame(where: Partial<BucketEntry>): Promise<Omit<BucketEntryWithFrame, 'frame'> & { frame?: Frame } | null>;
findByIdsWithFrames(ids: BucketEntry['id'][]): Promise<(Omit<BucketEntryWithFrame, 'frame'> & { frame?: Frame })[]>;
Expand Down
97 changes: 96 additions & 1 deletion lib/core/bucketEntries/usecase.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import lodash from 'lodash';

import log from '../../logger';
import { BucketsRepository } from '../buckets/Repository';
import { BucketEntriesRepository } from './Repository';
import { BucketNotFoundError, BucketForbiddenError, BucketEntryNotFoundError } from '../buckets/usecase';
Expand All @@ -16,6 +17,20 @@ import { User } from '../users/User';
import { Bucket } from '../buckets/Bucket';
import { FileStateRepository } from '../fileState/Repository';

/**
* Raised when a bucket turns out to hold entries backed by shards.
*
* Dropping those entries wholesale would strand their shards, mirrors and the
* bytes on the farmers with nothing left pointing at them.
*/
export class ShardBackedBucketError extends Error {
constructor() {
super('Bucket holds shard-backed entries and cannot be purged by this route');

Object.setPrototypeOf(this, ShardBackedBucketError.prototype);
}
}

export class BucketEntryVersionNotFoundError extends Error {
constructor() {
super('BucketEntryVersion not found');
Expand Down Expand Up @@ -179,8 +194,8 @@ export class BucketEntriesUsecase {
await this.bucketEntryShardsRepository.deleteByIds(bucketEntryShardsIds);
}

await this.bucketEntriesRepository.deleteByIds(fileIds);
await this.fileStateRepository.deleteByBucketEntryIds(fileIds);
await this.bucketEntriesRepository.deleteByIds(fileIds);
}

private async findBucketOwner(
Expand Down Expand Up @@ -254,4 +269,84 @@ export class BucketEntriesUsecase {
totalUsedSpaceBytes,
};
}

/**
* Removes a bucket and every entry in it.
*
* The entries this purges are metadata only: createEntry() writes a row with
* a bucket, a size and a version, and nothing else. Nothing downstream of a
* bucket entry exists for them, which is what makes a wholesale delete
* possible instead of walking each entry and its shards.
*
* A Drive bucket reaching this method would lose its files and strand their
* shards with nothing left to enumerate them by. The guards below read
* themselves, but note that the last of them - the delete only ever matching
* metadata-only entries - holds even if every check above it is wrong.
*
* BucketsUsecase.deleteBucketByIdAndUser is the legacy-stack equivalent; it
* drops the bucket document alone and does no usage accounting.
*/
async removeBucketAndEntries(
userUuid: User['uuid'],
bucketId: Bucket['id'],
bucketName: Bucket['name']
): Promise<UserSpaceSnapshot> {
const [user, bucket] = await Promise.all([
this.usersRepository.findByUuid(userUuid),
this.bucketsRepository.findOne({ id: bucketId }),
]);

if (!user) {
throw new UserNotFoundError(userUuid);
}

if (!bucket) {
throw new BucketNotFoundError();
}

if (bucket.userId !== userUuid) {
throw new BucketForbiddenError();
}

if (bucket.name !== bucketName) {
throw new BucketNotFoundError();
}

if (await this.bucketEntriesRepository.hasShardBackedEntriesByBucket(bucketId)) {
throw new ShardBackedBucketError();
}

const metadataOnlyBytes =
await this.bucketEntriesRepository.sumMetadataOnlyBytesByBucket(bucketId);

let totalUsedSpaceBytes = user.totalUsedSpaceBytes;

log.info(
`[removeBucketAndEntries] purging bucket - userUuid: ${userUuid}, bucketId: ${bucketId}, metadataOnlyBytes: ${metadataOnlyBytes}`
);

await this.bucketEntriesRepository.deleteMetadataOnlyByBucket(bucketId);

if (metadataOnlyBytes > 0) {
totalUsedSpaceBytes = await this.usersRepository.addTotalUsedSpaceBytes(
userUuid,
-metadataOnlyBytes
);
}

log.info(
`[removeBucketAndEntries] usage credited - userUuid: ${userUuid}, bucketId: ${bucketId}, releasedBytes: ${metadataOnlyBytes}, totalUsedSpaceBytes: ${totalUsedSpaceBytes}`
);

if (await this.bucketEntriesRepository.hasEntriesByBucket(bucketId)) {
throw new ShardBackedBucketError();
}

await this.bucketsRepository.removeByIdAndUser(bucketId, userUuid);

return {
maxSpaceBytes: user.maxSpaceBytes,
totalUsedSpaceBytes,
};
}
}
54 changes: 52 additions & 2 deletions lib/server/http/gateway/controller.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,8 @@ import { Request, Response } from 'express';
import { validate as uuidValidate } from 'uuid';
import { Logger } from 'winston';
import { EmailIsAlreadyInUseError, InvalidDataFormatError, UserAlreadyExistsError, UserNotFoundError, UserSpaceSnapshot, UsersUsecase } from '../../../core';
import { BucketEntriesUsecase } from '../../../core/bucketEntries/usecase';
import { BucketNotFoundError } from '../../../core/buckets/usecase';
import { BucketEntriesUsecase, ShardBackedBucketError } from '../../../core/bucketEntries/usecase';
import { BucketForbiddenError, BucketNotFoundError } from '../../../core/buckets/usecase';

import { GatewayUsecase } from '../../../core/gateway/Usecase';
import { EventBus, EventBusEvents, UserStorageChangedPayload } from '../../eventBus';
Expand Down Expand Up @@ -205,6 +205,56 @@ export class HTTPGatewayController {
}
}

async deleteUserBucket(
req: Request<{ uuid: string; id: string }, {}, {}, { name?: unknown }>,
res: Response<UserSpaceSnapshot | { message: string }>
) {
const { uuid, id } = req.params;
const { name } = req.query;

if (!uuid || !uuidValidate(uuid) || !id || !OBJECT_ID_PATTERN.test(id)) {
return res.status(400).send({ message: 'Invalid params' });
}

if (typeof name !== 'string' || name.length === 0) {
return res.status(400).send({ message: 'name is required' });
}

try {
const snapshot = await this.bucketEntriesUsecase.removeBucketAndEntries(uuid, id, name);

return res.status(200).send(snapshot);
} catch (err) {
if (err instanceof UserNotFoundError || err instanceof BucketNotFoundError) {
return res.status(404).send({ message: err.message });
}

if (err instanceof BucketForbiddenError) {
return res.status(403).send({ message: err.message });
}

if (err instanceof ShardBackedBucketError) {
this.logger.warn(
'[GATEWAY/DELETE_BUCKET] Refused to purge shard-backed bucket %s of user %s',
id,
uuid
);

return res.status(409).send({ message: err.message });
}

this.logger.error(
'[GATEWAY/DELETE_BUCKET] Error deleting bucket %s of user %s: %s. %s',
id,
uuid,
(err as Error).message,
(err as Error).stack || 'NO STACK'
);

return res.status(500).send({ message: 'Internal server error' });
}
}

async createBucketEntry(
req: Request<{ uuid: string; id: string }, {}, Partial<CreateBucketEntryBody>, {}>,
res: Response<CreateBucketEntryResponse | { message: string }>
Expand Down
1 change: 1 addition & 0 deletions lib/server/http/gateway/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ export const createGatewayHTTPRouter = (
router.patch('/users/:uuid', jwtMiddleware, controller.updateUserEmail.bind(controller));
router.put('/storage/users/:uuid', jwtMiddleware, controller.changeStorage.bind(controller));
router.post('/users/:uuid/buckets', jwtMiddleware, controller.createUserBucket.bind(controller));
router.delete('/users/:uuid/buckets/:id', jwtMiddleware, controller.deleteUserBucket.bind(controller));
router.post('/users/:uuid/buckets/:id/entries', jwtMiddleware, controller.createBucketEntry.bind(controller));
router.delete('/users/:uuid/buckets/:id/entries/:entryId', jwtMiddleware, controller.deleteBucketEntry.bind(controller));
router.delete('/storage/files', jwtMiddleware, controller.deleteFilesInBulk.bind(controller));
Expand Down
Loading
Loading