From bd4cf8b2850ea3a5dd2784ed51dc55cc97af9481 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Mon, 14 Sep 2026 18:03:37 +0700 Subject: [PATCH] refactor(fase1b): unify upload/download via use-case + proxy download - upload-file use-case is now the single save path with dedup policy (hash / bucket-key / none), partPrefix + signatureBuffer inputs, fileHash in UploadOutput, and owned temp-file cleanup; pass temp path (not open stream) to telegram service for testability - upload-controller (multipart + JSON) and web-api upload delegate to the use-case; JSON body via JsonUploadPayloadSchema; responses built from use-case output via buildUploadResponse - web-api download 302 redirect -> proxy stream (closes bot_token leak); shared CORS + sanitizeFilenameHeader - file-controller drops double file-info cache layer (pool caches) - bot handler uses DI singletons instead of direct construction - get-file RedirectRetrieval.redirectUrl marked deprecated, URL via shared builder; chunked-storage URL via buildTelegramFileUrl (no inline api.telegram.org/file/bot left in src/) --- src/application/dto/upload.ts | 26 +- src/application/use-cases/get-file.ts | 14 +- src/application/use-cases/upload-file.ts | 182 ++++++++------ .../telegram/chunked-storage.ts | 3 +- src/interfaces/bot/handler.ts | 13 +- .../http/controllers/file-controller.ts | 36 +-- .../http/controllers/upload-controller.ts | 238 ++++-------------- .../http/controllers/web-api-controller.ts | 171 +++++++------ 8 files changed, 292 insertions(+), 391 deletions(-) diff --git a/src/application/dto/upload.ts b/src/application/dto/upload.ts index fb70528..15761ea 100644 --- a/src/application/dto/upload.ts +++ b/src/application/dto/upload.ts @@ -1,3 +1,11 @@ +/** + * Deduplication policy for the upload file use case. + * - `hash`: look up an existing record by content SHA-256 first (default). + * - `bucket-key`: idempotent lookup by `(bucketId, s3Key)`; never re-uploads. + * - `none`: store unconditionally. + */ +export type UploadDedupPolicy = 'hash' | 'bucket-key' | 'none'; + /** * Input for the upload file use case. * Carries all metadata needed to persist an uploaded file, @@ -12,8 +20,11 @@ export interface UploadInput { fileName: string; /** MIME type detected from content inspection or request header */ mimeType: string; - /** High-level file category (e.g. "document", "photo", "video") */ - fileType: string; + /** + * High-level file category (e.g. "document", "photo", "video"). + * Optional — when omitted the use case derives it from MIME type + file name. + */ + fileType?: string; /** File size in bytes */ sizeBytes: number; /** Telegram user ID of the uploader; 0 when unknown or system */ @@ -22,6 +33,15 @@ export interface UploadInput { bucketId?: string | null; /** Object key within the bucket for S3-compatible storage; null when un-bucketed */ s3Key?: string | null; + /** Deduplication policy; defaults to `'hash'` when omitted */ + dedup?: UploadDedupPolicy; + /** Prefix for chunked part file names; defaults to `direct-` */ + partPrefix?: string; + /** + * First bytes of the file for magic-byte detection. + * When omitted the use case reads them from `tempPath`. + */ + signatureBuffer?: Buffer; } /** @@ -43,4 +63,6 @@ export interface UploadOutput { createdAt: Date; /** Public download URL */ downloadUrl: string; + /** SHA-256 hex digest of the stored content */ + fileHash: string | null; } diff --git a/src/application/use-cases/get-file.ts b/src/application/use-cases/get-file.ts index 4e1b81e..99ae571 100644 --- a/src/application/use-cases/get-file.ts +++ b/src/application/use-cases/get-file.ts @@ -1,6 +1,7 @@ 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. @@ -43,7 +44,13 @@ export interface RedirectRetrieval { type: 'redirect'; /** The resolved file entity. */ file: File; - /** Full Telegram CDN URL to redirect the client to. */ + /** + * 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; @@ -167,9 +174,10 @@ export function createGetFileUseCase(deps: GetFileUseCaseDeps) { return { type: 'archive-entry', file, archiveInfo, entryName: archiveEntryName }; } - // Regular file — resolve Telegram CDN path for a redirect + // 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 = `https://api.telegram.org/file/bot${fileInfo.bot_token}/${fileInfo.file_path}`; + const redirectUrl = buildTelegramFileUrl(fileInfo.file_path, fileInfo.bot_token); return { type: 'redirect', file, redirectUrl, fileInfo }; }; diff --git a/src/application/use-cases/upload-file.ts b/src/application/use-cases/upload-file.ts index caefc6c..031d8a2 100644 --- a/src/application/use-cases/upload-file.ts +++ b/src/application/use-cases/upload-file.ts @@ -1,11 +1,16 @@ -import { createReadStream } from 'node:fs'; import { open } from 'node:fs/promises'; import { nanoid } from 'nanoid'; +import type { File } from '../../domain/entities/file'; import { buildNewFile } from '../../domain/entities/file-factory'; import type { IFileRepository } from '../../domain/ports/file-repository'; import type { ITelegramService } from '../../domain/ports/telegram-service'; import type { ChunkedStorage } from '../../infrastructure/telegram/chunked-storage'; -import { checkFileSize, ensureExtension, getFileType } from '../../shared/utils/file'; +import { + checkFileSize, + cleanupTempFile, + ensureExtension, + getFileType, +} from '../../shared/utils/file'; import type { UploadInput, UploadOutput } from '../dto/upload'; /** Subset of application configuration consumed by the upload-file use case. */ @@ -51,11 +56,30 @@ const readSignatureBuffer = async (tempPath: string): Promise => { } }; +/** + * Maps a persisted file entity to the public upload output DTO. + * + * @param file - The persisted file record. + * @param baseUrl - Server base URL for the download link. + * @returns The public `UploadOutput` DTO. + */ +const toUploadOutput = (file: File, baseUrl: string): UploadOutput => ({ + publicId: file.publicId, + fileName: file.fileName, + mimeType: file.mimeType, + sizeBytes: file.sizeBytes, + fileType: file.fileType, + createdAt: file.createdAt instanceof Date ? file.createdAt : new Date(file.createdAt), + downloadUrl: `${baseUrl}/f/${file.publicId}`, + fileHash: file.fileHash, +}); + /** * 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). + * Single save path for all upload entry points (multipart, JSON, web-API): + * 1. Deduplication per policy — `hash` (by content SHA-256), `bucket-key` + * (idempotent by bucket + S3 key), or `none` (store unconditionally). * 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 (delegated to ChunkedStorage) or @@ -68,23 +92,26 @@ const readSignatureBuffer = async (tempPath: string): Promise => { */ 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}`, - }; + const dedup = input.dedup ?? 'hash'; + + // 1. Deduplication per policy + if (dedup === 'hash') { + const existing = await deps.fileRepo.findByHash(input.fileHash); + if (existing) { + return toUploadOutput(existing, deps.config.baseUrl); + } + } else if (dedup === 'bucket-key') { + if (!input.bucketId || !input.s3Key) { + throw new Error('bucket-key dedup requires bucketId and s3Key'); + } + const existing = await deps.fileRepo.findByBucketAndKey(input.bucketId, input.s3Key); + if (existing) { + return toUploadOutput(existing, deps.config.baseUrl); + } } - // 2. Read signature bytes for magic-byte-based extension detection - const signatureBuffer = await readSignatureBuffer(input.tempPath); + // 2. Signature bytes for magic-byte-based extension detection + const signatureBuffer = input.signatureBuffer ?? (await readSignatureBuffer(input.tempPath)); const { fileName: finalFileName, mimeType } = ensureExtension( input.fileName, @@ -93,73 +120,68 @@ export function createUploadFileUseCase(deps: UploadFileUseCaseDeps) { ); // 3. Determine Telegram file type and validate size - const fileType = getFileType(mimeType, finalFileName); + const fileTypeRaw = getFileType(mimeType, finalFileName); + const fileType = fileTypeRaw === 'application' ? 'document' : fileTypeRaw; if (!checkFileSize(input.sizeBytes, fileType)) { throw new Error(`File size exceeds ${fileType} limit`); } - // 4. Upload — chunked via ChunkedStorage for files above the threshold - if (input.sizeBytes > deps.config.telegramChunkSizeBytes) { - const uploadedFile = await deps.chunkedStorage.storeFileInTelegramChunks({ - tempPath: input.tempPath, - partFileNamePrefix: `direct-${input.fileHash.slice(0, 16)}`, - fileName: finalFileName, - mimeType, - sizeBytes: input.sizeBytes, - fileType, - uploaderId: input.uploaderId ?? 0, - bucketId: input.bucketId, - s3Key: input.s3Key, - }); + const partPrefix = input.partPrefix ?? `direct-${input.fileHash.slice(0, 16)}`; - return { - publicId: uploadedFile.publicId, - fileName: uploadedFile.fileName, - mimeType: uploadedFile.mimeType, - sizeBytes: uploadedFile.sizeBytes, - fileType: uploadedFile.fileType, - createdAt: uploadedFile.createdAt, - downloadUrl: `${deps.config.baseUrl}/f/${uploadedFile.publicId}`, - }; + // 4. Upload — chunked via ChunkedStorage for files above the threshold, + // single-message otherwise. The temp file is always cleaned up here so + // callers never need their own cleanup block. + try { + if (input.sizeBytes > deps.config.telegramChunkSizeBytes) { + const uploadedFile = await deps.chunkedStorage.storeFileInTelegramChunks({ + tempPath: input.tempPath, + partFileNamePrefix: partPrefix, + fileName: finalFileName, + mimeType, + sizeBytes: input.sizeBytes, + fileType, + uploaderId: input.uploaderId ?? 0, + bucketId: input.bucketId, + s3Key: input.s3Key, + }); + + return toUploadOutput(uploadedFile, deps.config.baseUrl); + } + + // 5. Single-message upload path. Pass the temp path (not an open + // stream) so the telegram service owns file I/O — test doubles that + // never touch disk keep working, and the use case stays stream-agnostic. + const forwardResult = await deps.telegramService.forwardToStorage( + input.tempPath, + finalFileName, + fileType, + ); + + const singlePublicId = nanoid(); + + const createdFile = await deps.fileRepo.create( + buildNewFile({ + publicId: singlePublicId, + telegramFileId: forwardResult.telegramFileId, + telegramFileUniqueId: forwardResult.telegramFileUniqueId, + storageChatId: deps.config.storageChatId, + storageMessageId: forwardResult.storageMessageId, + fileName: finalFileName, + mimeType, + sizeBytes: input.sizeBytes, + fileType, + storageBackend: 'telegram', + uploaderId: input.uploaderId, + fileHash: input.fileHash, + bucketId: input.bucketId, + s3Key: input.s3Key, + }), + ); + + return toUploadOutput(createdFile, deps.config.baseUrl); + } finally { + await cleanupTempFile(input.tempPath); } - - // 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( - buildNewFile({ - publicId: singlePublicId, - telegramFileId: forwardResult.telegramFileId, - telegramFileUniqueId: forwardResult.telegramFileUniqueId, - storageChatId: deps.config.storageChatId, - storageMessageId: forwardResult.storageMessageId, - fileName: finalFileName, - mimeType, - sizeBytes: input.sizeBytes, - fileType, - storageBackend: 'telegram', - uploaderId: input.uploaderId, - fileHash: input.fileHash, - bucketId: input.bucketId, - s3Key: input.s3Key, - }), - ); - - 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}`, - }; }; } diff --git a/src/infrastructure/telegram/chunked-storage.ts b/src/infrastructure/telegram/chunked-storage.ts index 54fe165..daacca7 100644 --- a/src/infrastructure/telegram/chunked-storage.ts +++ b/src/infrastructure/telegram/chunked-storage.ts @@ -12,6 +12,7 @@ import type { RangeParseResult } from '../../interfaces/s3/range'; import { type CompressionAlgorithm, maybeCompressChunk } from '../../shared/utils/compress'; import { computeHash } from '../../shared/utils/file'; import { asSafeChunkSize } from '../../shared/utils/validation'; +import { buildTelegramFileUrl } from './file-url'; /** * Metadata about a single uploaded chunk (part) stored in Telegram. @@ -237,7 +238,7 @@ export class ChunkedStorage { const fileInfo = await this.telegramService.getFileInfo(part.telegramFileId); return { part, - url: `https://api.telegram.org/file/bot${fileInfo.bot_token}/${fileInfo.file_path}`, + url: buildTelegramFileUrl(fileInfo.file_path, fileInfo.bot_token), }; }), ); diff --git a/src/interfaces/bot/handler.ts b/src/interfaces/bot/handler.ts index 69f5375..cd66620 100644 --- a/src/interfaces/bot/handler.ts +++ b/src/interfaces/bot/handler.ts @@ -4,8 +4,7 @@ import { buildNewFile } from '../../domain/entities/file-factory'; import type { IFileRepository } from '../../domain/ports/file-repository'; import type { ITelegramService } from '../../domain/ports/telegram-service'; import { config } from '../../env'; -import { DrizzleFileRepository } from '../../infrastructure/persistence/repositories/file-repository'; -import { botPool } from '../../infrastructure/telegram/bot-pool'; +import { fileRepository, telegramService } from '../../infrastructure/di'; import logger from '../../shared/logger/index'; import { checkFileSize, @@ -75,9 +74,9 @@ const replyWithDownloadUrl = async (ctx: BotContext, publicId: string): Promise< * * @param deps - Optional external dependencies for testing or DI override. * @param deps.telegramService - The Telegram service used to forward files to - * the storage channel. Defaults to the singleton BotPool instance. + * the storage channel. Defaults to the DI singleton. * @param deps.fileRepo - The file repository used for deduplication queries - * and persisting new file records. Defaults to a new DrizzleFileRepository. + * and persisting new file records. Defaults to the DI singleton. * @returns The launched Telegraf bot instance, suitable for graceful shutdown * via `bot.stop(signal)`. */ @@ -89,8 +88,8 @@ export async function startBot( fileRepo?: IFileRepository; } = {}, ): Promise> { - const telegramService = deps.telegramService ?? botPool; - const fileRepo = deps.fileRepo ?? new DrizzleFileRepository(); + const telegramSvc = deps.telegramService ?? telegramService; + const fileRepo = deps.fileRepo ?? fileRepository; try { const bot = new Telegraf(config.botTokens[0]); @@ -147,7 +146,7 @@ export async function startBot( return; } - const result = await telegramService.forwardToStorage(file_id, fileName, fileType); + const result = await telegramSvc.forwardToStorage(file_id, fileName, fileType); const publicId = nanoid(); await fileRepo.create( diff --git a/src/interfaces/http/controllers/file-controller.ts b/src/interfaces/http/controllers/file-controller.ts index ff2e7a4..91731ab 100644 --- a/src/interfaces/http/controllers/file-controller.ts +++ b/src/interfaces/http/controllers/file-controller.ts @@ -1,7 +1,5 @@ import { createReadStream } from 'node:fs'; import { nanoid } from 'nanoid'; -import type { TelegramFileInfo } from '../../../domain/ports/telegram-service'; -import { fileInfoCache } from '../../../infrastructure/cache/index'; import { chunkedStorage, fileRepository } from '../../../infrastructure/di'; import { botPool } from '../../../infrastructure/telegram/bot-pool'; import { buildTelegramFileUrl } from '../../../infrastructure/telegram/file-url'; @@ -21,33 +19,6 @@ type RequestWithParams = Request & { }; }; -/** - * Resolves Telegram file metadata for a given file ID, using the in-memory - * cache to avoid repeated API calls to Telegram. - * - * @param telegramFileId - The Telegram file identifier to resolve. - * @param publicId - The public file ID (used for logging). - * @returns The resolved Telegram file info. - */ -const getTelegramFileInfo = async ( - telegramFileId: string, - publicId: string, -): Promise => { - const cacheKey = `file_info_${telegramFileId}`; - const cached = fileInfoCache.get(cacheKey) as TelegramFileInfo | null; - - if (cached) { - logger.debug('File info from cache', { publicId, cacheKey }); - return cached; - } - - const fileInfo = await botPool.getFileInfo(telegramFileId); - fileInfoCache.set(cacheKey, fileInfo); - logger.debug('File info cached', { publicId, cacheKey }); - - return fileInfo; -}; - /** * Returns a JSON error response with the given status code and message. * @@ -65,7 +36,8 @@ const fail = (status: number, error: string): Response => Response.json({ error * - **chunked** files are streamed via the chunked-object response builder. * - **archive-entry** files are extracted from a Telegram-stored zip archive * and streamed as a single file. - * - **regular** files are redirected to the Telegram CDN URL (302). + * - **regular** files are proxied from the Telegram CDN (200 with the file + * body, so browser fetch/XHR playback never hits Telegram CORS). * * @param req - The incoming HTTP request with a `public_id` route parameter. * @returns A redirect or streaming response, or a JSON error. @@ -99,7 +71,7 @@ export const handleFileRedirect = async (req: RequestWithParams): Promise + createUploadFileUseCase({ + fileRepo: fileRepository, + telegramService, + chunkedStorage, + config: { + baseUrl: config.baseUrl, + telegramChunkSizeBytes: config.telegramChunkSizeBytes, + storageChatId: config.storageChatId, + compressChunkedUploads: config.compressChunkedUploads, + chunkCompressionMinSizeBytes: config.chunkCompressionMinSizeBytes, + }, + }); /** * Parses a base64-encoded file string, optionally stripping the data URI @@ -98,59 +94,13 @@ const rejectOversizedRequest = (req: Request): Response | null => { return null; }; -/** - * Streams a multipart `File` to a temporary file on disk while computing - * its SHA-256 hash and extracting the signature (first 16 bytes). - * - * Delegates to the shared {@link streamToTemp} utility. - * - * @param file - The multipart `File` object. - * @param maxSizeBytes - Maximum allowed file size; an error is thrown if - * the stream exceeds this limit. - * @returns A fully prepared upload descriptor with hash, size, and temp path. - * @throws {Error} When the file size exceeds `maxSizeBytes`. - */ -const streamFileToTemp = async (file: File, maxSizeBytes: number): Promise => { - const result = await streamToTemp(file.stream().getReader(), { maxSizeBytes }); - return result; -}; - -/** - * Writes an in-memory buffer to a temporary file on disk. - * - * Used for base64 JSON uploads where the decoded data is already in a Buffer. - * - * @param fileBuffer - The decoded file content. - * @param fileHash - Pre-computed SHA-256 hex digest. - * @returns A prepared upload descriptor. - */ -const writeBufferToTemp = async (fileBuffer: Buffer, fileHash: string): Promise => { - const tempPath = `/tmp/filedrop-${nanoid()}`; - try { - await Bun.write(tempPath, fileBuffer); - return { - tempPath, - fileHash, - sizeBytes: fileBuffer.byteLength, - signatureBuffer: fileBuffer.subarray(0, SIGNATURE_BYTES), - }; - } catch (error) { - await cleanupTempFile(tempPath); - throw error; - } -}; - /** * Handles a multipart/form-data file upload. * * Steps: * 1. Parse the multipart form and extract the file. * 2. Stream the file to a temp location, computing its hash. - * 3. Check for deduplication by content hash. - * 4. Determine the MIME type, file name, and Telegram file type. - * 5. Validate file size limits. - * 6. Upload to Telegram (chunked or single-message). - * 7. Return the upload response JSON. + * 3. Delegate to the upload use case (dedup `hash`) and return its response. * * @param req - The incoming HTTP request with a multipart body. * @returns A JSON response with the uploaded file metadata. @@ -170,70 +120,25 @@ const handleMultipartUpload = async (req: Request): Promise => { return Response.json({ error: 'File size exceeds upload limit' }, { status: 413 }); } - const prepared = await streamFileToTemp(file, config.maxRequestBodyBytes); - - const existingFile = await fileRepository.findByHash(prepared.fileHash); - if (existingFile) { - await cleanupTempFile(prepared.tempPath); - return Response.json(buildUploadResponse(existingFile, config.baseUrl), { status: 200 }); - } + const prepared = await streamToTemp(file.stream().getReader(), { + maxSizeBytes: config.maxRequestBodyBytes, + }); const rawMimeType = file.type || extractMimeType({}, req) || 'application/octet-stream'; - const { fileName: finalFileName, mimeType } = ensureExtension( + const output = await getUploadUseCase()({ + tempPath: prepared.tempPath, + fileHash: prepared.fileHash, fileName, - prepared.signatureBuffer, - rawMimeType, - ); - const fileType = getFileType(mimeType, finalFileName); + mimeType: rawMimeType, + sizeBytes: prepared.sizeBytes, + uploaderId: 0, + dedup: 'hash', + signatureBuffer: prepared.signatureBuffer, + }); - if (!checkFileSize(prepared.sizeBytes, fileType)) { - await cleanupTempFile(prepared.tempPath); - return Response.json({ error: `File size exceeds ${fileType} limit` }, { status: 400 }); - } - - if (prepared.sizeBytes > config.telegramChunkSizeBytes) { - const uploadedFile = await chunkedStorage.storeFileInTelegramChunks({ - tempPath: prepared.tempPath, - partFileNamePrefix: `direct-${prepared.fileHash?.slice(0, 16) || 'upload'}`, - fileName: finalFileName, - mimeType, - sizeBytes: prepared.sizeBytes, - fileType, - uploaderId: 0, - }); - await cleanupTempFile(prepared.tempPath); - return Response.json(buildUploadResponse(uploadedFile, config.baseUrl), { status: 200 }); - } - - // Single-message — direct to Telegram storage - const forwardResult = await telegramService.forwardToStorage( - Readable.from(Bun.file(prepared.tempPath).stream()), - finalFileName, - fileType, - ); - - const publicId = nanoid(); - - const createdFile = await fileRepository.create( - buildNewFile({ - publicId, - telegramFileId: forwardResult.telegramFileId, - telegramFileUniqueId: forwardResult.telegramFileUniqueId, - storageChatId: config.storageChatId, - storageMessageId: forwardResult.storageMessageId, - fileName: finalFileName, - mimeType, - sizeBytes: prepared.sizeBytes, - fileType, - storageBackend: 'telegram', - uploaderId: 0, - fileHash: prepared.fileHash, - }), - ); - - await cleanupTempFile(prepared.tempPath); - - return Response.json(buildUploadResponse(createdFile, config.baseUrl), { status: 200 }); + // The use case is the single source of truth for stored metadata — + // build the response directly from its output, not a synthetic record. + return Response.json(buildUploadResponse(output, config.baseUrl), { status: 200 }); } catch (error: unknown) { const message = getErrorMessage(error); logger.error('Multipart upload error', { error: message }); @@ -246,28 +151,23 @@ const handleMultipartUpload = async (req: Request): Promise => { * base64-encoded string. * * Steps: - * 1. Parse the JSON body and extract the base64 file data. - * 2. Decode and estimate the file size; reject if too large for JSON. - * 3. Write the decoded buffer to a temp file. - * 4. Check deduplication by content hash. - * 5. Determine MIME type, file name, and Telegram file type. - * 6. Validate file size limits. - * 7. Upload to Telegram (chunked or single-message). - * 8. Return the upload response JSON. + * 1. Parse and validate the JSON body (must include base64 `file`). + * 2. Decode, check size limits, and stage to a temp file. + * 3. Delegate to the upload use case (dedup `hash`) and return its response. * * @param req - The incoming HTTP request with a JSON body. * @returns A JSON response with the uploaded file metadata. */ const handleJSONUpload = async (req: Request): Promise => { try { - const { file, fileName = 'file' } = (await req.json()) as JsonUploadPayload; - - if (!file || typeof file !== 'string') { + const parsed = JsonUploadPayloadSchema.safeParse(await req.json()); + if (!parsed.success) { return Response.json( { error: 'Invalid JSON. Must include "file" (base64) and optional "fileName"' }, { status: 400 }, ); } + const { file, fileName } = parsed.data; const { base64Data, mimeType: rawMimeType } = parseBase64File(file); const estimatedSizeBytes = Math.floor((base64Data.length * 3) / 4); @@ -287,11 +187,6 @@ const handleJSONUpload = async (req: Request): Promise => { const fileBytes = Buffer.from(base64Data, 'base64'); const hash = computeHash(fileBytes); - const existingFile = await fileRepository.findByHash(hash); - if (existingFile) { - return Response.json(buildUploadResponse(existingFile, config.baseUrl), { status: 200 }); - } - const fileTypeRaw = getFileType(rawMimeType, fileName); const fileType = fileTypeRaw === 'application' ? 'document' : fileTypeRaw; @@ -301,51 +196,20 @@ const handleJSONUpload = async (req: Request): Promise => { return Response.json({ error: `File size exceeds ${fileType} limit` }, { status: 400 }); } - const prepared = await writeBufferToTemp(fileBytes, hash); + const tempPath = `/tmp/teleuploader-${nanoid()}`; + await Bun.write(tempPath, fileBytes); + const output = await getUploadUseCase()({ + tempPath, + fileHash: hash, + fileName: finalFileName, + mimeType, + sizeBytes: fileBytes.byteLength, + uploaderId: 0, + dedup: 'hash', + signatureBuffer: fileBytes.subarray(0, SIGNATURE_BYTES), + }); - if (prepared.sizeBytes > config.telegramChunkSizeBytes) { - const uploadedFile = await chunkedStorage.storeFileInTelegramChunks({ - tempPath: prepared.tempPath, - partFileNamePrefix: `direct-${prepared.fileHash?.slice(0, 16) || 'json'}`, - fileName: finalFileName, - mimeType, - sizeBytes: prepared.sizeBytes, - fileType, - uploaderId: 0, - }); - await cleanupTempFile(prepared.tempPath); - return Response.json(buildUploadResponse(uploadedFile, config.baseUrl), { status: 200 }); - } - - // Single-message — direct to Telegram storage - const forwardResult = await telegramService.forwardToStorage( - Readable.from(Bun.file(prepared.tempPath).stream()), - finalFileName, - fileType, - ); - - const publicId = nanoid(); - - const createdFile = await fileRepository.create( - buildNewFile({ - publicId, - telegramFileId: forwardResult.telegramFileId, - telegramFileUniqueId: forwardResult.telegramFileUniqueId, - storageChatId: config.storageChatId, - storageMessageId: forwardResult.storageMessageId, - fileName: finalFileName, - mimeType, - sizeBytes: prepared.sizeBytes, - fileType, - storageBackend: 'telegram', - uploaderId: 0, - fileHash: prepared.fileHash, - }), - ); - - await cleanupTempFile(prepared.tempPath); - - return Response.json(buildUploadResponse(createdFile, config.baseUrl), { status: 200 }); + return Response.json(buildUploadResponse(output, config.baseUrl), { status: 200 }); } catch (error: unknown) { const message = getErrorMessage(error); logger.error('JSON upload error', { error: message }); diff --git a/src/interfaces/http/controllers/web-api-controller.ts b/src/interfaces/http/controllers/web-api-controller.ts index 7af744a..76a4998 100644 --- a/src/interfaces/http/controllers/web-api-controller.ts +++ b/src/interfaces/http/controllers/web-api-controller.ts @@ -1,18 +1,32 @@ -import { createReadStream } from 'node:fs'; -import { nanoid } from 'nanoid'; -import { buildNewFile } from '../../../domain/entities/file-factory'; +import { createUploadFileUseCase } from '../../../application/use-cases/upload-file'; import { config } from '../../../env'; -import { bucketRepository, chunkedStorage, fileRepository } from '../../../infrastructure/di'; -import { botPool } from '../../../infrastructure/telegram/bot-pool'; -import logger from '../../../shared/logger/index'; import { - cleanupTempFile, - DEFAULT_FILE_TYPE, - ensureExtension, - getErrorMessage, -} from '../../../shared/utils/file'; + bucketRepository, + chunkedStorage, + fileRepository, + telegramService, +} from '../../../infrastructure/di'; +import { buildTelegramFileUrl } from '../../../infrastructure/telegram/file-url'; +import { sanitizeFilenameHeader } from '../../../shared/http/filename'; +import logger from '../../../shared/logger/index'; +import { getErrorMessage } from '../../../shared/utils/file'; import { streamToTemp } from '../../../shared/utils/temp-stream'; +/** Lazily built upload use case wired to the DI singletons. */ +const getUploadUseCase = () => + createUploadFileUseCase({ + fileRepo: fileRepository, + telegramService, + chunkedStorage, + config: { + baseUrl: config.baseUrl, + telegramChunkSizeBytes: config.telegramChunkSizeBytes, + storageChatId: config.storageChatId, + compressChunkedUploads: config.compressChunkedUploads, + chunkCompressionMinSizeBytes: config.chunkCompressionMinSizeBytes, + }, + }); + /** * Route parameters extracted from the URL path. */ @@ -170,73 +184,27 @@ export const handleUploadObjectV1 = async ( const key = (formData.get('key') as string) || file.name; const streamed = await streamToTemp(file.stream().getReader(), { prefix: '/tmp/filedrop-web-' }); - const { fileName: finalFileName, mimeType } = ensureExtension( - key.split('/').pop() || 'file', - streamed.signatureBuffer, - file.type || 'application/octet-stream', - ); - const partFileNamePrefix = `s3-${bucket.name}-${key.replace(/\//g, '_')}`; - - if (streamed.sizeBytes > config.telegramChunkSizeBytes) { - const uploadedFile = await chunkedStorage.storeFileInTelegramChunks({ - tempPath: streamed.tempPath, - partFileNamePrefix, - fileName: finalFileName, - mimeType, - sizeBytes: streamed.sizeBytes, - fileType: DEFAULT_FILE_TYPE, - uploaderId: 0, - bucketId: bucket.id, - s3Key: key, - }); - await cleanupTempFile(streamed.tempPath); - return json( - { - key, - size: streamed.sizeBytes, - etag: streamed.fileHash, - downloadUrl: `${config.baseUrl}/f/${uploadedFile.publicId}`, - }, - 201, - ); - } - - const forwardResult = await botPool.forwardToStorage( - createReadStream(streamed.tempPath), - partFileNamePrefix, - 'document', - ); - - const publicId = nanoid(); - - await fileRepository.create( - buildNewFile({ - publicId, - telegramFileId: forwardResult.telegramFileId, - telegramFileUniqueId: forwardResult.telegramFileUniqueId, - storageChatId: config.storageChatId, - storageMessageId: forwardResult.storageMessageId, - fileName: finalFileName, - mimeType, - sizeBytes: streamed.sizeBytes, - fileType: DEFAULT_FILE_TYPE, - uploaderId: 0, - fileHash: streamed.fileHash, - bucketId: bucket.id, - s3Key: key, - storageBackend: 'telegram', - }), - ); - - await cleanupTempFile(streamed.tempPath); + const output = await getUploadUseCase()({ + tempPath: streamed.tempPath, + fileHash: streamed.fileHash, + fileName: key.split('/').pop() || 'file', + mimeType: file.type || 'application/octet-stream', + sizeBytes: streamed.sizeBytes, + uploaderId: 0, + bucketId: bucket.id, + s3Key: key, + dedup: 'none', + partPrefix: `s3-${bucket.name}-${key.replace(/\//g, '_')}`, + signatureBuffer: streamed.signatureBuffer, + }); return json( { key, - size: streamed.sizeBytes, - etag: streamed.fileHash, - downloadUrl: `${config.baseUrl}/f/${publicId}`, + size: output.sizeBytes, + etag: output.fileHash, + downloadUrl: output.downloadUrl, }, 201, ); @@ -260,14 +228,15 @@ export const handleDeleteObjectV1 = async ( }; /** - * Downloads (or redirects to) an object from a bucket. + * Downloads (proxies) an object from a bucket. * * For chunked objects, builds a streaming response. For regular Telegram - * objects, issues a 302 redirect to the Telegram CDN URL. + * objects, proxies the Telegram CDN body (200 with the file body) so the + * bot token never leaks to clients via a redirect URL. * * @param _req - The incoming HTTP request (unused). * @param params - Route parameters containing the bucket name and object key. - * @returns A redirect or streaming response, or a JSON error. + * @returns A streaming response with the file body, or a JSON error. */ export const handleDownloadObjectV1 = async ( _req: Request, @@ -284,10 +253,54 @@ export const handleDownloadObjectV1 = async ( return chunkedStorage.createChunkedObjectResponse({ file, range, reqId: '' }); } - const fileInfo = await botPool.getFileInfo(file.telegramFileId); - const redirectUrl = `https://api.telegram.org/file/bot${fileInfo.bot_token}/${fileInfo.file_path}`; + const fileInfo = await telegramService.getFileInfo(file.telegramFileId); + const telegramUrl = buildTelegramFileUrl(fileInfo.file_path, fileInfo.bot_token); - return new Response(null, { status: 302, headers: { Location: redirectUrl } }); + const corsHeaders: Record = { + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Expose-Headers': 'Content-Disposition, Content-Length, Accept-Ranges', + 'Access-Control-Allow-Methods': 'GET, HEAD, OPTIONS', + 'Access-Control-Allow-Headers': 'Range, Content-Type', + Vary: 'Origin', + }; + + try { + const upstream = await fetch(telegramUrl); + if (!upstream.ok) { + logger.error('Web API object download failed', { + bucket: bucket.name, + key: params.key, + status: upstream.status, + }); + return Response.json( + { error: 'Upstream download failed' }, + { status: upstream.status === 404 ? 404 : 502, headers: corsHeaders }, + ); + } + + const upstreamHeaders = new Headers(upstream.headers); + return new Response(upstream.body, { + status: 200, + headers: { + 'Content-Type': + file.mimeType || upstreamHeaders.get('content-type') || 'application/octet-stream', + 'Content-Disposition': `attachment; filename="${sanitizeFilenameHeader(file.fileName)}"`, + 'Content-Length': String(file.sizeBytes ?? 0), + 'Cache-Control': 'public, max-age=300', + ...corsHeaders, + }, + }); + } catch (error: unknown) { + logger.error('Web API object proxy error', { + bucket: bucket.name, + key: params.key, + error: getErrorMessage(error), + }); + return Response.json( + { error: 'Upstream download failed' }, + { status: 502, headers: corsHeaders }, + ); + } }; /**