refactor(plan): delete dead use-cases s3-object, manage-bucket, multipart-upload, get-file
Verified zero importers in src/ and test/ (only a doc comment mentions manage-bucket). Live use-cases: authenticate (auth controller) and upload-file (upload + web-api controllers); S3 handlers call repositories directly. Build + biome + full test:unit (28 files) green.
This commit is contained in:
@@ -1,184 +0,0 @@
|
|||||||
import type { File } from '../../domain/entities/file';
|
|
||||||
import type { IFileRepository } from '../../domain/ports/file-repository';
|
|
||||||
import type { ITelegramService, TelegramFileInfo } from '../../domain/ports/telegram-service';
|
|
||||||
import { buildTelegramFileUrl } from '../../infrastructure/telegram/file-url';
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result type for a simple file-info lookup.
|
|
||||||
*/
|
|
||||||
export interface FileInfoResult {
|
|
||||||
/** Whether the file was found. */
|
|
||||||
found: true;
|
|
||||||
/** Public unique identifier. */
|
|
||||||
publicId: string;
|
|
||||||
/** Original file name. */
|
|
||||||
fileName: string;
|
|
||||||
/** MIME type. */
|
|
||||||
mimeType: string;
|
|
||||||
/** File size in bytes. */
|
|
||||||
sizeBytes: number;
|
|
||||||
/** Telegram file type (document, photo, video, etc.). */
|
|
||||||
fileType: string;
|
|
||||||
/** ISO-8601 creation timestamp. */
|
|
||||||
createdAt: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result type for a file-not-found lookup.
|
|
||||||
*/
|
|
||||||
export interface FileNotFoundResult {
|
|
||||||
/** Always `false` for a not-found result. */
|
|
||||||
found: false;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Discriminated union of all possible file-info lookup outcomes.
|
|
||||||
*/
|
|
||||||
export type GetFileInfoResult = FileInfoResult | FileNotFoundResult;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Describes a redirect-based file retrieval.
|
|
||||||
*/
|
|
||||||
export interface RedirectRetrieval {
|
|
||||||
/** Discriminant. */
|
|
||||||
type: 'redirect';
|
|
||||||
/** The resolved file entity. */
|
|
||||||
file: File;
|
|
||||||
/**
|
|
||||||
* Full Telegram CDN URL to redirect the client to.
|
|
||||||
*
|
|
||||||
* @deprecated No controller redirects anymore — downloads are proxied
|
|
||||||
* server-side via `buildTelegramFileUrl` so the bot token never reaches
|
|
||||||
* clients. Kept for API compatibility; do not expose to clients.
|
|
||||||
*/
|
|
||||||
redirectUrl: string;
|
|
||||||
/** Cached Telegram file metadata. */
|
|
||||||
fileInfo: TelegramFileInfo;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Describes a chunked file retrieval that needs a multi-part response.
|
|
||||||
*/
|
|
||||||
export interface ChunkedRetrieval {
|
|
||||||
/** Discriminant. */
|
|
||||||
type: 'chunked';
|
|
||||||
/** The resolved file entity. */
|
|
||||||
file: File;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Describes an archive-entry file retrieval.
|
|
||||||
*/
|
|
||||||
export interface ArchiveEntryRetrieval {
|
|
||||||
/** Discriminant. */
|
|
||||||
type: 'archive-entry';
|
|
||||||
/** The resolved file entity. */
|
|
||||||
file: File;
|
|
||||||
/** Telegram file metadata for the archive container. */
|
|
||||||
archiveInfo: TelegramFileInfo;
|
|
||||||
/** Name of the entry within the archive. */
|
|
||||||
entryName: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Discriminated union of all possible file retrieval outcomes.
|
|
||||||
*/
|
|
||||||
export type FileRetrievalResult = RedirectRetrieval | ChunkedRetrieval | ArchiveEntryRetrieval;
|
|
||||||
|
|
||||||
/** Subset of application configuration consumed by the get-file use case. */
|
|
||||||
export interface GetFileConfig {
|
|
||||||
/** Server base URL (used in constructing archive download URLs). */
|
|
||||||
baseUrl: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Dependencies required by the get-file use case factory. */
|
|
||||||
export interface GetFileUseCaseDeps {
|
|
||||||
/** File repository for looking up file records. */
|
|
||||||
fileRepo: IFileRepository;
|
|
||||||
/** Telegram service for resolving file identifiers to download paths. */
|
|
||||||
telegramService: ITelegramService;
|
|
||||||
/** Application configuration subset. */
|
|
||||||
config: GetFileConfig;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a factory function for the get-file info use case.
|
|
||||||
*
|
|
||||||
* Looks up a file by its public identifier and returns its metadata.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting a public ID and returning file info.
|
|
||||||
*/
|
|
||||||
export function createGetFileInfoUseCase(deps: Pick<GetFileUseCaseDeps, 'fileRepo'>) {
|
|
||||||
return async (publicId: string): Promise<GetFileInfoResult> => {
|
|
||||||
const file = await deps.fileRepo.findByPublicId(publicId);
|
|
||||||
if (!file) {
|
|
||||||
return { found: false };
|
|
||||||
}
|
|
||||||
|
|
||||||
return {
|
|
||||||
found: true,
|
|
||||||
publicId: file.publicId,
|
|
||||||
fileName: file.fileName,
|
|
||||||
mimeType: file.mimeType,
|
|
||||||
sizeBytes: file.sizeBytes,
|
|
||||||
fileType: file.fileType,
|
|
||||||
createdAt: formatCreatedAtForInfo(file.createdAt),
|
|
||||||
};
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Formats a date-like value into an ISO-8601 string.
|
|
||||||
*
|
|
||||||
* @param date - A Date instance, date string, or numeric timestamp.
|
|
||||||
* @returns The ISO-8601 string.
|
|
||||||
*/
|
|
||||||
const formatCreatedAtForInfo = (date: Date | string | number): string => {
|
|
||||||
if (date instanceof Date) return date.toISOString();
|
|
||||||
return new Date(date).toISOString();
|
|
||||||
};
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a factory function for the get-file retrieval use case.
|
|
||||||
*
|
|
||||||
* Determines how a file should be delivered to the client:
|
|
||||||
* - **redirect**: For regular (non-chunked, non-archive) files — returns a
|
|
||||||
* Telegram CDN redirect URL.
|
|
||||||
* - **chunked**: For files stored across multiple Telegram parts — returns
|
|
||||||
* the file entity so the caller can build a multi-part streaming response.
|
|
||||||
* - **archive-entry**: For files stored inside a Telegram archive (zip) —
|
|
||||||
* returns the archive's Telegram metadata and the entry name so the caller
|
|
||||||
* can extract and stream the entry.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting a public ID and returning a retrieval result.
|
|
||||||
*/
|
|
||||||
export function createGetFileUseCase(deps: GetFileUseCaseDeps) {
|
|
||||||
return async (publicId: string): Promise<FileRetrievalResult | null> => {
|
|
||||||
const file = await deps.fileRepo.findByPublicId(publicId);
|
|
||||||
if (!file) {
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Chunked file — return the entity for multi-part response building
|
|
||||||
if (file.storageBackend === 'chunked') {
|
|
||||||
return { type: 'chunked', file };
|
|
||||||
}
|
|
||||||
|
|
||||||
// Archive entry — resolve the archive's Telegram location
|
|
||||||
const archiveEntryName = file.archiveEntryName;
|
|
||||||
if (archiveEntryName) {
|
|
||||||
const archiveFileId = file.archiveTelegramFileId || file.telegramFileId;
|
|
||||||
const archiveInfo = await deps.telegramService.getFileInfo(archiveFileId);
|
|
||||||
return { type: 'archive-entry', file, archiveInfo, entryName: archiveEntryName };
|
|
||||||
}
|
|
||||||
|
|
||||||
// Regular file — resolve Telegram CDN path (server-side fetch only;
|
|
||||||
// the URL embeds the bot token and must never be exposed to clients).
|
|
||||||
const fileInfo = await deps.telegramService.getFileInfo(file.telegramFileId);
|
|
||||||
const redirectUrl = buildTelegramFileUrl(fileInfo.file_path, fileInfo.bot_token);
|
|
||||||
|
|
||||||
return { type: 'redirect', file, redirectUrl, fileInfo };
|
|
||||||
};
|
|
||||||
}
|
|
||||||
@@ -1,156 +0,0 @@
|
|||||||
import type { Bucket } from '../../domain/entities/bucket';
|
|
||||||
import type { IBucketRepository } from '../../domain/ports/bucket-repository';
|
|
||||||
import type { IFileRepository } from '../../domain/ports/file-repository';
|
|
||||||
|
|
||||||
/**
|
|
||||||
* S3 bucket name validation regex.
|
|
||||||
*
|
|
||||||
* Bucket names must be 3-63 characters, start/end with a lowercase letter or
|
|
||||||
* digit, and contain only lowercase letters, digits, dots, and hyphens.
|
|
||||||
*/
|
|
||||||
const BUCKET_NAME_REGEX = /^[a-z0-9][a-z0-9.-]{1,61}[a-z0-9]$/;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Error type for bucket-level application errors that carry an S3-compatible
|
|
||||||
* error code and an HTTP status suggestion.
|
|
||||||
*/
|
|
||||||
export class BucketError extends Error {
|
|
||||||
/** S3-compatible error code (e.g. "NoSuchBucket", "BucketAlreadyExists"). */
|
|
||||||
readonly code: string;
|
|
||||||
/** Suggested HTTP status code. */
|
|
||||||
readonly status: number;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @param code - The S3 error code.
|
|
||||||
* @param message - Human-readable error description.
|
|
||||||
* @param status - Suggested HTTP status.
|
|
||||||
*/
|
|
||||||
constructor(code: string, message: string, status: number) {
|
|
||||||
super(message);
|
|
||||||
this.name = 'BucketError';
|
|
||||||
this.code = code;
|
|
||||||
this.status = status;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Describes a bucket entry returned by the list-buckets use case,
|
|
||||||
* enriched with the current count of non-deleted objects.
|
|
||||||
*/
|
|
||||||
export interface BucketWithCount {
|
|
||||||
/** The bucket domain entity. */
|
|
||||||
bucket: Bucket;
|
|
||||||
/** Number of non-deleted objects in the bucket. */
|
|
||||||
objectCount: number;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Dependencies required by the manage-bucket use case factories. */
|
|
||||||
export interface ManageBucketDeps {
|
|
||||||
/** Bucket repository for CRUD operations. */
|
|
||||||
bucketRepo: IBucketRepository;
|
|
||||||
/** File repository for counting and checking objects within buckets. */
|
|
||||||
fileRepo: IFileRepository;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that lists all buckets together with their object counts.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function that returns a list of buckets with counts.
|
|
||||||
*/
|
|
||||||
export function createListBucketsUseCase(deps: ManageBucketDeps) {
|
|
||||||
return async (): Promise<BucketWithCount[]> => {
|
|
||||||
const buckets = await deps.bucketRepo.list();
|
|
||||||
const results = await Promise.all(
|
|
||||||
buckets.map(async (bucket) => ({
|
|
||||||
bucket,
|
|
||||||
objectCount: await deps.fileRepo.countByBucket(bucket.id),
|
|
||||||
})),
|
|
||||||
);
|
|
||||||
return results;
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that retrieves a single bucket by its name.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting a bucket name and returning the
|
|
||||||
* bucket, or `null` when not found.
|
|
||||||
*/
|
|
||||||
export function createGetBucketUseCase(deps: ManageBucketDeps) {
|
|
||||||
return async (name: string): Promise<Bucket | null> => {
|
|
||||||
return deps.bucketRepo.findByName(name);
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that creates a new bucket.
|
|
||||||
*
|
|
||||||
* Validates the bucket name format and checks for duplicates before
|
|
||||||
* persisting.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting a bucket name and returning the
|
|
||||||
* newly created bucket.
|
|
||||||
* @throws {BucketError} When the name is invalid or the bucket already exists.
|
|
||||||
*/
|
|
||||||
export function createCreateBucketUseCase(deps: ManageBucketDeps) {
|
|
||||||
return async (name: string): Promise<Bucket> => {
|
|
||||||
if (!BUCKET_NAME_REGEX.test(name)) {
|
|
||||||
throw new BucketError('InvalidBucketName', 'The specified bucket is not valid.', 400);
|
|
||||||
}
|
|
||||||
|
|
||||||
const existing = await deps.bucketRepo.findByName(name);
|
|
||||||
if (existing) {
|
|
||||||
throw new BucketError(
|
|
||||||
'BucketAlreadyExists',
|
|
||||||
'The requested bucket name is not available.',
|
|
||||||
409,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
return deps.bucketRepo.create(name);
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that deletes a bucket.
|
|
||||||
*
|
|
||||||
* Ensures the bucket exists and is empty (no non-deleted objects) before
|
|
||||||
* proceeding with deletion.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting a bucket name. Returns `true` when
|
|
||||||
* the bucket was deleted, throws when the bucket is missing or
|
|
||||||
* not empty.
|
|
||||||
* @throws {BucketError} When the bucket does not exist or is not empty.
|
|
||||||
*/
|
|
||||||
export function createDeleteBucketUseCase(deps: ManageBucketDeps) {
|
|
||||||
return async (name: string): Promise<boolean> => {
|
|
||||||
const bucket = await deps.bucketRepo.findByName(name);
|
|
||||||
if (!bucket) {
|
|
||||||
throw new BucketError('NoSuchBucket', 'The specified bucket does not exist.', 404);
|
|
||||||
}
|
|
||||||
|
|
||||||
const objectCount = await deps.fileRepo.countByBucket(bucket.id);
|
|
||||||
if (objectCount > 0) {
|
|
||||||
throw new BucketError('BucketNotEmpty', 'The bucket you tried to delete is not empty.', 409);
|
|
||||||
}
|
|
||||||
|
|
||||||
return deps.bucketRepo.delete(name);
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that checks whether a bucket exists.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting a bucket name and returning `true`
|
|
||||||
* when the bucket exists.
|
|
||||||
*/
|
|
||||||
export function createBucketExistsUseCase(deps: ManageBucketDeps) {
|
|
||||||
return async (name: string): Promise<boolean> => {
|
|
||||||
return deps.bucketRepo.exists(name);
|
|
||||||
};
|
|
||||||
}
|
|
||||||
@@ -1,368 +0,0 @@
|
|||||||
import { nanoid } from 'nanoid';
|
|
||||||
import { buildNewFile } from '../../domain/entities/file-factory';
|
|
||||||
import type { MultipartUpload } from '../../domain/entities/multipart';
|
|
||||||
import type { IBucketRepository } from '../../domain/ports/bucket-repository';
|
|
||||||
import type { IFileRepository } from '../../domain/ports/file-repository';
|
|
||||||
import type { IMultipartRepository } from '../../domain/ports/multipart-repository';
|
|
||||||
import type { ITelegramService } from '../../domain/ports/telegram-service';
|
|
||||||
import { computeHash, DEFAULT_FILE_TYPE } from '../../shared/utils/file';
|
|
||||||
|
|
||||||
// ─── Types ──────────────────────────────────────────────────────────
|
|
||||||
|
|
||||||
/**
|
|
||||||
* A single part reference as submitted in a complete-multipart-upload request.
|
|
||||||
*/
|
|
||||||
export interface CompletePartInput {
|
|
||||||
/** 1-based part number. */
|
|
||||||
partNumber: number;
|
|
||||||
/** ETag returned when the part was uploaded. */
|
|
||||||
etag: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result of initiating a multipart upload.
|
|
||||||
*/
|
|
||||||
export interface InitiateMultipartResult {
|
|
||||||
/** The generated upload identifier (nanoid). */
|
|
||||||
uploadId: string;
|
|
||||||
/** The bucket name. */
|
|
||||||
bucket: string;
|
|
||||||
/** The S3 object key. */
|
|
||||||
key: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result of uploading a single part.
|
|
||||||
*/
|
|
||||||
export interface UploadPartResult {
|
|
||||||
/** ETag of the uploaded part (SHA-256 hex digest). */
|
|
||||||
etag: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result of completing a multipart upload.
|
|
||||||
*/
|
|
||||||
export interface CompleteMultipartResult {
|
|
||||||
/** Public-facing unique identifier of the created file record. */
|
|
||||||
publicId: string;
|
|
||||||
/** The S3 location URL of the completed object. */
|
|
||||||
location: string;
|
|
||||||
/** Combined ETag (all part etags joined by hyphens). */
|
|
||||||
etag: string;
|
|
||||||
/** Total object size in bytes. */
|
|
||||||
sizeBytes: number;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Summary of a single part within a multipart upload used in listing results.
|
|
||||||
*/
|
|
||||||
export interface PartSummary {
|
|
||||||
/** 1-based part number. */
|
|
||||||
partNumber: number;
|
|
||||||
/** ETag of the part content. */
|
|
||||||
etag: string;
|
|
||||||
/** Part size in bytes. */
|
|
||||||
sizeBytes: number;
|
|
||||||
/** ISO-8601 timestamp when the part was stored. */
|
|
||||||
createdAt: Date;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result of listing multipart uploads within a bucket.
|
|
||||||
*/
|
|
||||||
export interface ListMultipartUploadsResult {
|
|
||||||
/** Array of in-progress upload summaries. */
|
|
||||||
uploads: MultipartUpload[];
|
|
||||||
/** Whether more results are available. */
|
|
||||||
isTruncated: boolean;
|
|
||||||
/** Marker for the next page, or null when not truncated. */
|
|
||||||
nextKeyMarker: string | null;
|
|
||||||
}
|
|
||||||
|
|
||||||
// ─── Config ─────────────────────────────────────────────────────────
|
|
||||||
|
|
||||||
/** Subset of application configuration consumed by the multipart use cases. */
|
|
||||||
export interface MultipartConfig {
|
|
||||||
/** Maximum chunk size in bytes for Telegram uploads (part size limit). */
|
|
||||||
telegramChunkSizeBytes: number;
|
|
||||||
/** Telegram chat ID where part data is stored. */
|
|
||||||
storageChatId: number;
|
|
||||||
/** Server base URL for constructing location URLs. */
|
|
||||||
baseUrl: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Dependencies required by the multipart upload use case factories. */
|
|
||||||
export interface MultipartDeps {
|
|
||||||
/** Bucket repository for bucket lookups. */
|
|
||||||
bucketRepo: IBucketRepository;
|
|
||||||
/** File repository for creating the final file record on completion. */
|
|
||||||
fileRepo: IFileRepository;
|
|
||||||
/** Multipart repository for managing upload sessions and parts. */
|
|
||||||
multipartRepo: IMultipartRepository;
|
|
||||||
/** Telegram service for forwarding part data to storage. */
|
|
||||||
telegramService: ITelegramService;
|
|
||||||
/** Application configuration subset. */
|
|
||||||
config: MultipartConfig;
|
|
||||||
}
|
|
||||||
|
|
||||||
// ─── Use Case Factories ─────────────────────────────────────────────
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that initiates an S3 multipart upload.
|
|
||||||
*
|
|
||||||
* Validates the bucket exists and creates a new multipart upload session.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting bucket name and object key, returning
|
|
||||||
* the upload initiation result, or `null` when the bucket is not found.
|
|
||||||
*/
|
|
||||||
export function createInitiateMultipartUploadUseCase(deps: MultipartDeps) {
|
|
||||||
return async (bucketName: string, key: string): Promise<InitiateMultipartResult | null> => {
|
|
||||||
const bucket = await deps.bucketRepo.findByName(bucketName);
|
|
||||||
if (!bucket) return null;
|
|
||||||
|
|
||||||
const uploadId = await deps.multipartRepo.create(bucket.id, key, 's3');
|
|
||||||
|
|
||||||
return { uploadId, bucket: bucketName, key };
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that uploads a single part of a multipart upload.
|
|
||||||
*
|
|
||||||
* Validates the part number range (1-10000), checks the upload session exists
|
|
||||||
* and matches the expected key, checks part size against the configured limit,
|
|
||||||
* forwards the part data to Telegram storage, and persists the part record.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting upload details and part data, returning
|
|
||||||
* the part ETag, or `null` when the upload session is not found.
|
|
||||||
*/
|
|
||||||
export function createUploadPartUseCase(deps: MultipartDeps) {
|
|
||||||
return async (input: {
|
|
||||||
/** Bucket name for the multipart upload. */
|
|
||||||
bucketName: string;
|
|
||||||
/** S3 object key for the multipart upload. */
|
|
||||||
key: string;
|
|
||||||
/** Upload identifier returned by initiate. */
|
|
||||||
uploadId: string;
|
|
||||||
/** 1-based part number (1-10000). */
|
|
||||||
partNumber: number;
|
|
||||||
/** Raw part data. */
|
|
||||||
body: Buffer;
|
|
||||||
}): Promise<UploadPartResult | null> => {
|
|
||||||
if (input.partNumber < 1 || input.partNumber > 10000) {
|
|
||||||
throw new MultipartError(
|
|
||||||
'InvalidArgument',
|
|
||||||
'Part number must be an integer between 1 and 10000',
|
|
||||||
400,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
const multipart = await deps.multipartRepo.findById(input.uploadId);
|
|
||||||
if (!multipart || multipart.s3Key !== input.key) {
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (input.body.byteLength > deps.config.telegramChunkSizeBytes) {
|
|
||||||
throw new MultipartError(
|
|
||||||
'EntityTooLarge',
|
|
||||||
`Your proposed upload part size (${input.body.byteLength} bytes) exceeds the maximum allowed part size (${deps.config.telegramChunkSizeBytes} bytes) for this storage backend. Use smaller part sizes.`,
|
|
||||||
400,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
const forwardResult = await deps.telegramService.forwardToStorage(
|
|
||||||
input.body,
|
|
||||||
`mp-${input.uploadId}-part-${input.partNumber}`,
|
|
||||||
'document',
|
|
||||||
);
|
|
||||||
|
|
||||||
const etag = computeHash(input.body);
|
|
||||||
|
|
||||||
await deps.multipartRepo.insertPart({
|
|
||||||
uploadId: input.uploadId,
|
|
||||||
partNumber: input.partNumber,
|
|
||||||
telegramFileId: forwardResult.telegramFileId,
|
|
||||||
telegramFileUniqueId: forwardResult.telegramFileUniqueId,
|
|
||||||
storageMessageId: forwardResult.storageMessageId,
|
|
||||||
sizeBytes: input.body.byteLength,
|
|
||||||
etag,
|
|
||||||
});
|
|
||||||
|
|
||||||
return { etag: `"${etag}"` };
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Error type for multipart-level application errors.
|
|
||||||
*/
|
|
||||||
export class MultipartError extends Error {
|
|
||||||
/** S3-compatible error code. */
|
|
||||||
readonly code: string;
|
|
||||||
/** Suggested HTTP status code. */
|
|
||||||
readonly status: number;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* @param code - The S3 error code.
|
|
||||||
* @param message - Human-readable error description.
|
|
||||||
* @param status - Suggested HTTP status.
|
|
||||||
*/
|
|
||||||
constructor(code: string, message: string, status: number) {
|
|
||||||
super(message);
|
|
||||||
this.name = 'MultipartError';
|
|
||||||
this.code = code;
|
|
||||||
this.status = status;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that completes an S3 multipart upload.
|
|
||||||
*
|
|
||||||
* Validates the submitted part list (all parts must be present and in ascending
|
|
||||||
* order), creates the final file record referencing the first part's Telegram
|
|
||||||
* data, marks the upload session as completed, and returns the combined result.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting upload details and submitted parts,
|
|
||||||
* returning the completion result, or `null` when the upload session
|
|
||||||
* is not found.
|
|
||||||
*/
|
|
||||||
export function createCompleteMultipartUploadUseCase(deps: MultipartDeps) {
|
|
||||||
return async (input: {
|
|
||||||
/** Bucket name for the multipart upload. */
|
|
||||||
bucketName: string;
|
|
||||||
/** S3 object key for the multipart upload. */
|
|
||||||
key: string;
|
|
||||||
/** Upload identifier. */
|
|
||||||
uploadId: string;
|
|
||||||
/** Parts submitted by the client (in ascending part number order). */
|
|
||||||
parts: CompletePartInput[];
|
|
||||||
}): Promise<CompleteMultipartResult | null> => {
|
|
||||||
const multipart = await deps.multipartRepo.findById(input.uploadId);
|
|
||||||
if (!multipart) return null;
|
|
||||||
|
|
||||||
const storedParts = await deps.multipartRepo.listParts(input.uploadId);
|
|
||||||
|
|
||||||
// Validate ascending part order
|
|
||||||
const partNumbers = input.parts.map((p) => p.partNumber);
|
|
||||||
if (partNumbers.length > 1 && partNumbers.some((n, i) => i > 0 && n <= partNumbers[i - 1])) {
|
|
||||||
throw new MultipartError(
|
|
||||||
'InvalidPartOrder',
|
|
||||||
'The list of parts was not in ascending order.',
|
|
||||||
400,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
// Validate part count matches
|
|
||||||
if (input.parts.length !== storedParts.length) {
|
|
||||||
throw new MultipartError(
|
|
||||||
'InvalidPart',
|
|
||||||
'One or more specified parts could not be found.',
|
|
||||||
400,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
const totalSize = storedParts.reduce((sum, p) => sum + p.sizeBytes, 0);
|
|
||||||
const firstPart = storedParts[0];
|
|
||||||
if (!firstPart) {
|
|
||||||
throw new MultipartError('InternalError', 'Multipart object has no parts.', 500);
|
|
||||||
}
|
|
||||||
|
|
||||||
const publicId = nanoid();
|
|
||||||
|
|
||||||
await deps.fileRepo.create(
|
|
||||||
buildNewFile({
|
|
||||||
publicId,
|
|
||||||
telegramFileId: firstPart.telegramFileId,
|
|
||||||
telegramFileUniqueId: firstPart.telegramFileUniqueId,
|
|
||||||
storageChatId: deps.config.storageChatId,
|
|
||||||
storageMessageId: firstPart.storageMessageId,
|
|
||||||
fileName: input.key.split('/').pop() || 'file',
|
|
||||||
mimeType: 'application/octet-stream',
|
|
||||||
sizeBytes: totalSize,
|
|
||||||
fileType: DEFAULT_FILE_TYPE,
|
|
||||||
storageBackend: 'telegram',
|
|
||||||
bucketId: multipart.bucketId,
|
|
||||||
s3Key: input.key,
|
|
||||||
multipartUploadId: input.uploadId,
|
|
||||||
}),
|
|
||||||
);
|
|
||||||
|
|
||||||
await deps.multipartRepo.complete(input.uploadId);
|
|
||||||
|
|
||||||
const location = `${deps.config.baseUrl}/${input.bucketName}/${input.key}`;
|
|
||||||
const combinedEtag = storedParts.map((p) => p.etag).join('-');
|
|
||||||
|
|
||||||
return {
|
|
||||||
publicId,
|
|
||||||
location,
|
|
||||||
etag: combinedEtag,
|
|
||||||
sizeBytes: totalSize,
|
|
||||||
};
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that aborts an S3 multipart upload.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting an upload identifier, returning `true`
|
|
||||||
* when the upload was aborted, or `null` when the upload session
|
|
||||||
* is not found.
|
|
||||||
*/
|
|
||||||
export function createAbortMultipartUploadUseCase(deps: MultipartDeps) {
|
|
||||||
return async (uploadId: string): Promise<boolean | null> => {
|
|
||||||
const multipart = await deps.multipartRepo.findById(uploadId);
|
|
||||||
if (!multipart) return null;
|
|
||||||
|
|
||||||
await deps.multipartRepo.abort(uploadId);
|
|
||||||
return true;
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that lists in-progress multipart uploads within a bucket.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting query parameters and returning the
|
|
||||||
* listing result, or `null` when the bucket is not found.
|
|
||||||
*/
|
|
||||||
export function createListMultipartUploadsUseCase(deps: MultipartDeps) {
|
|
||||||
return async (input: {
|
|
||||||
/** Bucket name to list uploads from. */
|
|
||||||
bucketName: string;
|
|
||||||
/** Maximum number of uploads to return (clamped 1-1000). */
|
|
||||||
maxUploads: number;
|
|
||||||
/** Return only uploads whose S3 key is strictly greater than this, or null. */
|
|
||||||
keyMarker: string | null;
|
|
||||||
}): Promise<ListMultipartUploadsResult | null> => {
|
|
||||||
const bucket = await deps.bucketRepo.findByName(input.bucketName);
|
|
||||||
if (!bucket) return null;
|
|
||||||
|
|
||||||
return deps.multipartRepo.listByBucket(bucket.id, input.maxUploads, input.keyMarker);
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Creates a use case that lists parts of a specific multipart upload.
|
|
||||||
*
|
|
||||||
* @param deps - The injected dependencies.
|
|
||||||
* @returns An async function accepting an upload identifier and returning the
|
|
||||||
* list of parts, or `null` when the upload session is not found.
|
|
||||||
*/
|
|
||||||
export function createListPartsUseCase(deps: MultipartDeps) {
|
|
||||||
return async (uploadId: string): Promise<PartSummary[] | null> => {
|
|
||||||
const multipart = await deps.multipartRepo.findById(uploadId);
|
|
||||||
if (!multipart) return null;
|
|
||||||
|
|
||||||
const parts = await deps.multipartRepo.listParts(uploadId);
|
|
||||||
|
|
||||||
return parts.map((p) => ({
|
|
||||||
partNumber: p.partNumber,
|
|
||||||
etag: p.etag,
|
|
||||||
sizeBytes: p.sizeBytes,
|
|
||||||
createdAt: p.createdAt,
|
|
||||||
}));
|
|
||||||
};
|
|
||||||
}
|
|
||||||
@@ -1,165 +0,0 @@
|
|||||||
import type { File } from '../../domain/entities/file';
|
|
||||||
import type { IBucketRepository } from '../../domain/ports/bucket-repository';
|
|
||||||
import type { IFilePartRepository } from '../../domain/ports/file-part-repository';
|
|
||||||
import type { IFileRepository } from '../../domain/ports/file-repository';
|
|
||||||
import type { IMultipartRepository } from '../../domain/ports/multipart-repository';
|
|
||||||
import type { ITelegramService, TelegramFileInfo } from '../../domain/ports/telegram-service';
|
|
||||||
import type { CompressionAlgorithm } from '../../shared/utils/compress';
|
|
||||||
|
|
||||||
// ─── Types ──────────────────────────────────────────────────────────
|
|
||||||
|
|
||||||
/**
|
|
||||||
* A single part source for building a multi-part streaming response.
|
|
||||||
* Each part corresponds to a Telegram-stored file chunk.
|
|
||||||
*/
|
|
||||||
export interface ObjectPartSource {
|
|
||||||
/** Telegram file identifier for retrieving this part. */
|
|
||||||
telegramFileId: string;
|
|
||||||
/** Telegram CDN URL for downloading this part. */
|
|
||||||
telegramUrl: string;
|
|
||||||
/** Original size of this part in bytes. */
|
|
||||||
sizeBytes: number;
|
|
||||||
/** 1-based part number within the object. */
|
|
||||||
partNumber: number;
|
|
||||||
/** Stored (post-compression) size in bytes, when applicable. */
|
|
||||||
storedSizeBytes?: number;
|
|
||||||
/** Compression algorithm applied, or null if uncompressed. */
|
|
||||||
compressionAlgorithm?: CompressionAlgorithm;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* A regular (direct) S3 object resolved to a Telegram CDN URL.
|
|
||||||
*/
|
|
||||||
export interface DirectObjectResult {
|
|
||||||
/** Discriminant. */
|
|
||||||
type: 'direct';
|
|
||||||
/** The resolved file entity. */
|
|
||||||
file: File;
|
|
||||||
/** Full Telegram CDN URL for downloading the object. */
|
|
||||||
telegramUrl: string;
|
|
||||||
/** Telegram file metadata. */
|
|
||||||
fileInfo: TelegramFileInfo;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* A chunked S3 object stored across multiple Telegram file parts.
|
|
||||||
*/
|
|
||||||
export interface ChunkedObjectResult {
|
|
||||||
/** Discriminant. */
|
|
||||||
type: 'chunked';
|
|
||||||
/** The resolved file entity. */
|
|
||||||
file: File;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* An S3 object assembled from a completed multipart upload.
|
|
||||||
*/
|
|
||||||
export interface MultipartObjectResult {
|
|
||||||
/** Discriminant. */
|
|
||||||
type: 'multipart';
|
|
||||||
/** The resolved file entity. */
|
|
||||||
file: File;
|
|
||||||
/** Resolved part sources with Telegram CDN URLs. */
|
|
||||||
parts: ObjectPartSource[];
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Discriminated union of all possible S3 get-object outcomes.
|
|
||||||
*/
|
|
||||||
export type GetObjectResult = DirectObjectResult | ChunkedObjectResult | MultipartObjectResult;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result of an S3 put-object operation.
|
|
||||||
*/
|
|
||||||
export interface PutObjectResult {
|
|
||||||
/** SHA-256 hex digest of the object content. */
|
|
||||||
etag: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result of an S3 copy-object operation.
|
|
||||||
*/
|
|
||||||
export interface CopyObjectResult {
|
|
||||||
/** SHA-256 hex digest of the source object content. */
|
|
||||||
etag: string;
|
|
||||||
/** ISO-8601 timestamp of the copy operation. */
|
|
||||||
lastModified: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* A single S3 object as returned in listing results.
|
|
||||||
*/
|
|
||||||
export interface ListObjectEntry {
|
|
||||||
/** The object key (full path within the bucket). */
|
|
||||||
key: string;
|
|
||||||
/** Object size in bytes. */
|
|
||||||
sizeBytes: number;
|
|
||||||
/** SHA-256 hex digest or fallback identifier. */
|
|
||||||
etag: string;
|
|
||||||
/** ISO-8601 timestamp of last modification. */
|
|
||||||
lastModified: string;
|
|
||||||
/** MIME type of the stored object. */
|
|
||||||
mimeType: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Result of an S3 list-objects operation (both V1 and V2).
|
|
||||||
*/
|
|
||||||
export interface ListObjectsResult {
|
|
||||||
/** Array of object summaries. */
|
|
||||||
objects: ListObjectEntry[];
|
|
||||||
/** Common prefixes when a delimiter was used. */
|
|
||||||
prefixes: string[];
|
|
||||||
/** Whether more results are available. */
|
|
||||||
isTruncated: boolean;
|
|
||||||
/** The last key in the returned page, for use as the next marker. */
|
|
||||||
nextMarker: string | null;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Head-object metadata.
|
|
||||||
*/
|
|
||||||
export interface HeadObjectMetadata {
|
|
||||||
/** MIME type of the object. */
|
|
||||||
contentType: string;
|
|
||||||
/** Object size in bytes. */
|
|
||||||
contentLength: number;
|
|
||||||
/** Entity tag (SHA-256 hex digest or fallback). */
|
|
||||||
etag: string;
|
|
||||||
/** ISO-8601 timestamp of last modification. */
|
|
||||||
lastModified: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
// ─── Config ─────────────────────────────────────────────────────────
|
|
||||||
|
|
||||||
/** Subset of application configuration consumed by the s3-object use cases. */
|
|
||||||
export interface S3ObjectConfig {
|
|
||||||
/** Maximum chunk size in bytes for Telegram chunked uploads. */
|
|
||||||
telegramChunkSizeBytes: number;
|
|
||||||
/** Whether gzip compression is enabled for chunked uploads. */
|
|
||||||
compressChunkedUploads: boolean;
|
|
||||||
/** Minimum chunk size in bytes below which compression is skipped. */
|
|
||||||
chunkCompressionMinSizeBytes: number;
|
|
||||||
/** Telegram chat ID where file parts are stored. */
|
|
||||||
storageChatId: number;
|
|
||||||
/** Server base URL for constructing download links. */
|
|
||||||
baseUrl: string;
|
|
||||||
/** Whether to proxy S3 GET requests through the server. */
|
|
||||||
proxyS3Get: boolean;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Dependencies required by the s3-object use case factories. */
|
|
||||||
export interface S3ObjectDeps {
|
|
||||||
/** Bucket repository for bucket lookups. */
|
|
||||||
bucketRepo: IBucketRepository;
|
|
||||||
/** File repository for object CRUD operations. */
|
|
||||||
fileRepo: IFileRepository;
|
|
||||||
/** File-part repository for chunked upload part records. */
|
|
||||||
filePartRepo: IFilePartRepository;
|
|
||||||
/** Multipart repository for resolving multipart-upload objects. */
|
|
||||||
multipartRepo: IMultipartRepository;
|
|
||||||
/** Telegram service for uploading and resolving file metadata. */
|
|
||||||
telegramService: ITelegramService;
|
|
||||||
/** Application configuration subset. */
|
|
||||||
config: S3ObjectConfig;
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user