From e2de8c424536b0c24c329e1dd048123a0efbf321 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 28 Jul 2026 17:59:49 +0700 Subject: [PATCH] feat: create application use cases (upload, get-file, auth) Create upload-file.ts use case with factory pattern supporting dedup, file type detection, size validation, and chunked/single storage strategies. Create get-file.ts use case supporting redirect, chunked, and archive-entry retrieval strategies. Create authenticate.ts use case with login, logout, and me operations. All use cases use dependency injection and return typed DTOs. Co-Authored-By: Claude Opus 5 (1M context) --- src/application/use-cases/authenticate.ts | 108 +++++++ src/application/use-cases/get-file.ts | 176 +++++++++++ src/application/use-cases/upload-file.ts | 360 ++++++++++++++++++++++ 3 files changed, 644 insertions(+) create mode 100644 src/application/use-cases/authenticate.ts create mode 100644 src/application/use-cases/get-file.ts create mode 100644 src/application/use-cases/upload-file.ts diff --git a/src/application/use-cases/authenticate.ts b/src/application/use-cases/authenticate.ts new file mode 100644 index 0000000..e9a3097 --- /dev/null +++ b/src/application/use-cases/authenticate.ts @@ -0,0 +1,108 @@ +import { timingSafeEqual } from 'node:crypto'; +import type { LoginInput, LoginResponse, LogoutResponse, UserInfoResponse, AuthSession } from '../dto/auth'; + +/** Subset of application configuration consumed by the authenticate use case. */ +export interface AuthUseCaseConfig { + /** Admin API token used to authenticate login requests. */ + adminApiToken: string; + /** Name of the session cookie. */ + sessionCookieName: string; + /** Session lifetime in milliseconds. */ + sessionMaxAgeMs: number; +} + +/** Dependencies required by the authenticate use case factory. */ +export interface AuthenticateUseCaseDeps { + /** Application configuration subset. */ + config: AuthUseCaseConfig; +} + +/** + * Performs a constant-time string comparison to prevent timing attacks. + * + * @param left - The first string to compare. + * @param right - The second string to compare. + * @returns `true` if the strings are equal, `false` otherwise. + */ +const timingSafeCompare = (left: string, right: string): boolean => { + const leftBuffer = Buffer.from(left); + const rightBuffer = Buffer.from(right); + + if (leftBuffer.length !== rightBuffer.length) { + return false; + } + + return timingSafeEqual(leftBuffer, rightBuffer); +}; + +/** + * Checks whether authentication is enabled based on the configured token. + * + * @param adminApiToken - The admin API token value. + * @returns `true` if the token is non-empty (auth is enabled). + */ +const isAuthEnabled = (adminApiToken: string): boolean => adminApiToken.length > 0; + +/** + * Creates a factory function for the login use case. + * + * Validates the provided admin API token and returns session metadata on + * success. The caller (controller/adapter) is responsible for translating + * the result into an HTTP response (e.g. setting a session cookie). + * + * @param deps - The injected dependencies. + * @returns An async function accepting login input and returning a login response. + */ +export function createLoginUseCase(deps: AuthenticateUseCaseDeps) { + return async (input: LoginInput): Promise => { + if (!isAuthEnabled(deps.config.adminApiToken)) { + return { username: 'admin' }; + } + + if (!timingSafeCompare(input.token, deps.config.adminApiToken)) { + throw new Error('Invalid token'); + } + + return { username: 'admin' }; + }; +} + +/** + * Creates a factory function for the logout use case. + * + * Always succeeds — the caller is responsible for clearing the session cookie. + * + * @returns An async function returning a logout response. + */ +export function createLogoutUseCase() { + return async (): Promise => { + return { success: true }; + }; +} + +/** + * Creates a factory function for the current-user (me) use case. + * + * Accepts an already-parsed auth session (from cookie or bearer token) and + * returns the user info response. The caller (controller/adapter) is + * responsible for extracting the session from the raw HTTP request. + * + * @param deps - The injected dependencies. + * @returns An async function accepting an optional session and returning user info. + */ +export function createMeUseCase(deps: AuthenticateUseCaseDeps) { + return async (session: AuthSession | null): Promise => { + if (!isAuthEnabled(deps.config.adminApiToken)) { + return null; + } + + if (!session) { + return null; + } + + return { + username: session.username, + expiresAt: session.expiresAt?.toISOString() ?? null, + }; + }; +} \ No newline at end of file diff --git a/src/application/use-cases/get-file.ts b/src/application/use-cases/get-file.ts new file mode 100644 index 0000000..69bf8ab --- /dev/null +++ b/src/application/use-cases/get-file.ts @@ -0,0 +1,176 @@ +import type { File } from '../../domain/entities/file'; +import type { IFileRepository } from '../../domain/ports/file-repository'; +import type { ITelegramService, TelegramFileInfo } from '../../domain/ports/telegram-service'; + +/** + * 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. */ + 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) { + return async (publicId: string): Promise => { + 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 => { + 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 for a redirect + const fileInfo = await deps.telegramService.getFileInfo(file.telegramFileId); + const redirectUrl = `https://api.telegram.org/file/bot${fileInfo.bot_token}/${fileInfo.file_path}`; + + return { type: 'redirect', file, redirectUrl, fileInfo }; + }; +} \ No newline at end of file diff --git a/src/application/use-cases/upload-file.ts b/src/application/use-cases/upload-file.ts new file mode 100644 index 0000000..3dfd30c --- /dev/null +++ b/src/application/use-cases/upload-file.ts @@ -0,0 +1,360 @@ +import { randomUUID } from 'node:crypto'; +import { open } from 'node:fs/promises'; +import { createReadStream } from 'node:fs'; +import { gzipSync } from 'node:zlib'; +import { nanoid } from 'nanoid'; +import type { NewFilePart } from '../../domain/entities/file-part'; +import type { IFilePartRepository } from '../../domain/ports/file-part-repository'; +import type { IFileRepository } from '../../domain/ports/file-repository'; +import type { ITelegramService } from '../../domain/ports/telegram-service'; +import type { UploadInput, UploadOutput } from '../dto/upload'; +import { getFileType, checkFileSize, ensureExtension, computeHash, formatCreatedAt } from '../../shared/utils/file'; + +/** Compression algorithm string literal used in chunked storage. */ +type ChunkCompressionAlgorithm = 'gzip' | null; + +/** Metadata for a single uploaded chunk/part. */ +interface UploadedPart { + /** 1-based part number. */ + partNumber: number; + /** Telegram file identifier for this part. */ + telegramFileId: string; + /** Telegram unique file identifier (stable across bot tokens). */ + telegramFileUniqueId: string; + /** Message ID within the storage chat. */ + storageMessageId: number; + /** Original size of the chunk in bytes before compression. */ + sizeBytes: number; + /** Stored (post-compression) size in bytes. */ + storedSizeBytes: number; + /** Compression algorithm applied, or null if uncompressed. */ + compressionAlgorithm: ChunkCompressionAlgorithm; + /** SHA-256 hash of the original chunk content. */ + etag: string; +} + +/** Result of uploading a file in multiple Telegram chunks. */ +interface ChunkedUploadResult { + /** Metadata for each uploaded part. */ + parts: UploadedPart[]; + /** SHA-256 hex digest of the complete file content. */ + fileHash: string; + /** Total file size in bytes (sum of all original chunks). */ + totalSizeBytes: number; +} + +/** Subset of application configuration consumed by the upload-file use case. */ +export interface UploadFileConfig { + /** Server base URL for constructing download links. */ + baseUrl: string; + /** Maximum chunk size in bytes for Telegram chunked uploads. */ + telegramChunkSizeBytes: number; + /** Telegram chat ID where file parts are stored. */ + storageChatId: number; + /** Whether to attempt gzip compression on each chunk. */ + compressChunkedUploads: boolean; + /** Minimum chunk size in bytes below which compression is skipped. */ + chunkCompressionMinSizeBytes: number; +} + +/** Dependencies required by the upload-file use case factory. */ +export interface UploadFileUseCaseDeps { + /** File repository for CRUD operations on file records. */ + fileRepo: IFileRepository; + /** File-part repository for chunked file metadata. */ + filePartRepo: IFilePartRepository; + /** Telegram service for forwarding file content to storage. */ + telegramService: ITelegramService; + /** Application configuration subset. */ + config: UploadFileConfig; +} + +/** + * Validates the configured chunk size and returns it as a safe integer. + * + * @param chunkSizeBytes - The configured chunk size in bytes. + * @returns The same value if it is a positive safe integer. + */ +const asSafeChunkSize = (chunkSizeBytes: number): number => { + if (!Number.isSafeInteger(chunkSizeBytes) || chunkSizeBytes <= 0) { + throw new Error('Invalid Telegram chunk size'); + } + return chunkSizeBytes; +}; + +/** + * Optionally gzip-compresses a chunk if compression is enabled and the chunk + * is large enough to benefit from it. + * + * @param chunk - The raw chunk buffer. + * @param compress - Whether compression is enabled. + * @param compressionMinSizeBytes - Minimum chunk size to attempt compression. + * @returns The (possibly compressed) buffer and the algorithm used. + */ +const maybeCompressChunk = ( + chunk: Buffer, + compress: boolean, + compressionMinSizeBytes: number, +): { bytes: Buffer; compressionAlgorithm: ChunkCompressionAlgorithm } => { + if (!compress || chunk.byteLength < compressionMinSizeBytes) { + return { bytes: chunk, compressionAlgorithm: null }; + } + + const gzipped = gzipSync(chunk); + if (gzipped.byteLength >= chunk.byteLength) { + return { bytes: chunk, compressionAlgorithm: null }; + } + + return { bytes: gzipped, compressionAlgorithm: 'gzip' }; +}; + +/** + * Reads the first 16 bytes from a file on disk for magic-byte detection. + * + * @param tempPath - Absolute path to the temporary file. + * @returns A buffer containing up to 16 bytes. + */ +const readSignatureBuffer = async (tempPath: string): Promise => { + const handle = await open(tempPath, 'r'); + try { + const buf = Buffer.alloc(16); + const { bytesRead } = await handle.read(buf, 0, 16, 0); + return buf.subarray(0, bytesRead); + } finally { + await handle.close(); + } +}; + +/** + * Reads a file from disk in chunks, forwards each chunk to Telegram storage, + * and returns metadata for all uploaded parts together with the total file + * hash. + * + * @param tempPath - Absolute path to the temporary file on disk. + * @param partFileNamePrefix - Prefix used for each chunk's file name in Telegram. + * @param chunkSizeBytes - Maximum size of each chunk in bytes. + * @param compress - Whether gzip compression is enabled. + * @param compressionMinSizeBytes - Minimum chunk size to attempt compression. + * @param telegramService - The Telegram service to forward each chunk. + * @returns The aggregated chunked upload result. + */ +const uploadFileInTelegramChunks = async ( + tempPath: string, + partFileNamePrefix: string, + chunkSizeBytes: number, + compress: boolean, + compressionMinSizeBytes: number, + telegramService: ITelegramService, +): Promise => { + const safeChunkSize = asSafeChunkSize(chunkSizeBytes); + const hasher = new Bun.CryptoHasher('sha256'); + const parts: UploadedPart[] = []; + let totalSizeBytes = 0; + let partNumber = 0; + + const stream = createReadStream(tempPath, { highWaterMark: safeChunkSize }); + + for await (const data of stream) { + const chunk = Buffer.isBuffer(data) ? data : Buffer.from(data as Uint8Array); + if (chunk.byteLength === 0) continue; + + partNumber += 1; + totalSizeBytes += chunk.byteLength; + hasher.update(chunk); + + const { bytes, compressionAlgorithm } = maybeCompressChunk( + chunk, + compress, + compressionMinSizeBytes, + ); + + const forwardResult = await telegramService.forwardToStorage( + bytes, + `${partFileNamePrefix}.part-${partNumber}`, + 'document', + ); + + parts.push({ + partNumber, + telegramFileId: forwardResult.telegramFileId, + telegramFileUniqueId: forwardResult.telegramFileUniqueId, + storageMessageId: forwardResult.storageMessageId, + sizeBytes: chunk.byteLength, + storedSizeBytes: bytes.byteLength, + compressionAlgorithm, + etag: computeHash(chunk), + }); + } + + return { + parts, + fileHash: hasher.digest('hex'), + totalSizeBytes, + }; +}; + +/** + * Creates a factory function for the upload-file use case. + * + * The returned use case: + * 1. Checks for an existing file with the same SHA-256 hash (deduplication). + * 2. Normalises the file name and MIME type based on magic bytes. + * 3. Validates the file size against Telegram type-specific limits. + * 4. Chooses a storage strategy — chunked (for files exceeding the chunk + * threshold) or single-message upload. + * 5. Persists the file record (and, for chunked uploads, part records). + * 6. Builds and returns the public `UploadOutput` DTO. + * + * @param deps - The injected dependencies. + * @returns An async function accepting `UploadInput` and returning `UploadOutput`. + */ +export function createUploadFileUseCase(deps: UploadFileUseCaseDeps) { + return async (input: UploadInput): Promise => { + // 1. Check deduplication by content hash + const existing = await deps.fileRepo.findByHash(input.fileHash); + if (existing) { + return { + publicId: existing.publicId, + fileName: existing.fileName, + mimeType: existing.mimeType, + sizeBytes: existing.sizeBytes, + fileType: existing.fileType, + createdAt: existing.createdAt instanceof Date ? existing.createdAt : new Date(existing.createdAt), + downloadUrl: `${deps.config.baseUrl}/f/${existing.publicId}`, + }; + } + + // 2. Read signature bytes for magic-byte-based extension detection + const signatureBuffer = await readSignatureBuffer(input.tempPath); + + const { fileName: finalFileName, mimeType } = ensureExtension( + input.fileName, + signatureBuffer, + input.mimeType, + ); + + // 3. Determine Telegram file type and validate size + const fileType = getFileType(mimeType, finalFileName); + + if (!checkFileSize(input.sizeBytes, fileType)) { + throw new Error(`File size exceeds ${fileType} limit`); + } + + // 4. Upload — chunked for files above the threshold, single otherwise + if (input.sizeBytes > deps.config.telegramChunkSizeBytes) { + // Chunked upload path + const chunkResult = await uploadFileInTelegramChunks( + input.tempPath, + `direct-${input.fileHash.slice(0, 16)}`, + deps.config.telegramChunkSizeBytes, + deps.config.compressChunkedUploads, + deps.config.chunkCompressionMinSizeBytes, + deps.telegramService, + ); + + const firstPart = chunkResult.parts[0]; + if (!firstPart) { + throw new Error('Chunked upload produced no parts'); + } + + const fileId = randomUUID(); + const publicId = nanoid(); + + const newFile = await deps.fileRepo.create({ + publicId, + telegramFileId: firstPart.telegramFileId, + telegramFileUniqueId: firstPart.telegramFileUniqueId, + storageChatId: deps.config.storageChatId, + storageMessageId: firstPart.storageMessageId, + fileName: finalFileName, + mimeType, + sizeBytes: chunkResult.totalSizeBytes, + fileType, + uploaderId: input.uploaderId ?? 0, + fileHash: chunkResult.fileHash, + archiveTelegramFileId: null, + archiveStorageMessageId: null, + archiveFileName: null, + archiveEntryName: null, + archiveMimeType: null, + archiveSizeBytes: null, + bucketId: input.bucketId ?? null, + s3Key: input.s3Key ?? null, + storageBackend: 'chunked', + isDeleted: false, + multipartUploadId: null, + partCount: chunkResult.parts.length, + }); + + const fileParts: NewFilePart[] = chunkResult.parts.map((part) => ({ + fileId, + partNumber: part.partNumber, + telegramFileId: part.telegramFileId, + telegramFileUniqueId: part.telegramFileUniqueId, + storageChatId: deps.config.storageChatId, + storageMessageId: part.storageMessageId, + sizeBytes: part.sizeBytes, + storedSizeBytes: part.storedSizeBytes, + compressionAlgorithm: part.compressionAlgorithm, + etag: part.etag, + })); + + await deps.filePartRepo.insert(fileParts); + + return { + publicId: newFile.publicId, + fileName: newFile.fileName, + mimeType: newFile.mimeType, + sizeBytes: newFile.sizeBytes, + fileType: newFile.fileType, + createdAt: newFile.createdAt, + downloadUrl: `${deps.config.baseUrl}/f/${newFile.publicId}`, + }; + } + + // 5. Single-message upload path + const forwardResult = await deps.telegramService.forwardToStorage( + createReadStream(input.tempPath), + finalFileName, + fileType, + ); + + const singlePublicId = nanoid(); + + const createdFile = await deps.fileRepo.create({ + publicId: singlePublicId, + telegramFileId: forwardResult.telegramFileId, + telegramFileUniqueId: forwardResult.telegramFileUniqueId, + storageChatId: deps.config.storageChatId, + storageMessageId: forwardResult.storageMessageId, + fileName: finalFileName, + mimeType, + sizeBytes: input.sizeBytes, + fileType, + uploaderId: input.uploaderId ?? 0, + fileHash: input.fileHash, + archiveTelegramFileId: null, + archiveStorageMessageId: null, + archiveFileName: null, + archiveEntryName: null, + archiveMimeType: null, + archiveSizeBytes: null, + bucketId: input.bucketId ?? null, + s3Key: input.s3Key ?? null, + storageBackend: 'telegram', + isDeleted: false, + multipartUploadId: null, + partCount: null, + }); + + return { + publicId: createdFile.publicId, + fileName: createdFile.fileName, + mimeType: createdFile.mimeType, + sizeBytes: createdFile.sizeBytes, + fileType: createdFile.fileType, + createdAt: createdFile.createdAt, + downloadUrl: `${deps.config.baseUrl}/f/${createdFile.publicId}`, + }; + }; +} \ No newline at end of file