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/)
This commit is contained in:
@@ -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-<hash16>` */
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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 };
|
||||
};
|
||||
|
||||
@@ -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<Buffer> => {
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
* 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<Buffer> => {
|
||||
*/
|
||||
export function createUploadFileUseCase(deps: UploadFileUseCaseDeps) {
|
||||
return async (input: UploadInput): Promise<UploadOutput> => {
|
||||
// 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}`,
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
};
|
||||
}),
|
||||
);
|
||||
|
||||
@@ -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<Telegraf<Context>> {
|
||||
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(
|
||||
|
||||
@@ -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<TelegramFileInfo> => {
|
||||
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<Respon
|
||||
const archiveEntryName = file.archiveEntryName;
|
||||
if (archiveEntryName) {
|
||||
const archiveFileId = file.archiveTelegramFileId || file.telegramFileId;
|
||||
const archiveInfo = await getTelegramFileInfo(archiveFileId, publicId);
|
||||
const archiveInfo = await botPool.getFileInfo(archiveFileId);
|
||||
const archiveResponse = await fetch(
|
||||
buildTelegramFileUrl(archiveInfo.file_path, archiveInfo.bot_token),
|
||||
);
|
||||
@@ -143,7 +115,7 @@ export const handleFileRedirect = async (req: RequestWithParams): Promise<Respon
|
||||
});
|
||||
}
|
||||
|
||||
const fileInfo = await getTelegramFileInfo(file.telegramFileId, publicId);
|
||||
const fileInfo = await botPool.getFileInfo(file.telegramFileId);
|
||||
const telegramUrl = buildTelegramFileUrl(fileInfo.file_path, fileInfo.bot_token);
|
||||
|
||||
// CORS headers shared across all delivery modes — public file CDN.
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
import { Readable } from 'node:stream';
|
||||
import { nanoid } from 'nanoid';
|
||||
import { buildNewFile } from '../../../domain/entities/file-factory';
|
||||
import { createUploadFileUseCase } from '../../../application/use-cases/upload-file';
|
||||
import { config } from '../../../env';
|
||||
import { chunkedStorage, fileRepository, telegramService } from '../../../infrastructure/di';
|
||||
import logger from '../../../shared/logger/index';
|
||||
@@ -8,7 +7,6 @@ import { metricsCollector } from '../../../shared/metrics/index';
|
||||
import {
|
||||
buildUploadResponse,
|
||||
checkFileSize,
|
||||
cleanupTempFile,
|
||||
computeHash,
|
||||
ensureExtension,
|
||||
extractMimeType,
|
||||
@@ -16,14 +14,7 @@ import {
|
||||
getFileType,
|
||||
} from '../../../shared/utils/file';
|
||||
import { streamToTemp } from '../../../shared/utils/temp-stream';
|
||||
|
||||
/** Prepared upload metadata before submission to storage. */
|
||||
interface PreparedUpload {
|
||||
tempPath: string;
|
||||
fileHash: string;
|
||||
sizeBytes: number;
|
||||
signatureBuffer: Buffer;
|
||||
}
|
||||
import { JsonUploadPayloadSchema } from '../../../shared/validation/schemas';
|
||||
|
||||
/**
|
||||
* Maximum allowed size (in bytes) for a base64 JSON upload.
|
||||
@@ -35,15 +26,20 @@ const JSON_UPLOAD_LIMIT_BYTES = 50 * 1024 * 1024;
|
||||
/** Number of leading bytes read for magic-byte / signature detection. */
|
||||
const SIGNATURE_BYTES = 16;
|
||||
|
||||
/**
|
||||
* Payload structure accepted by the JSON upload endpoint.
|
||||
*/
|
||||
interface JsonUploadPayload {
|
||||
/** Base64-encoded file data (optionally with a data URI prefix). */
|
||||
file?: unknown;
|
||||
/** Optional file name. */
|
||||
fileName?: string;
|
||||
}
|
||||
/** 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,
|
||||
},
|
||||
});
|
||||
|
||||
/**
|
||||
* 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<PreparedUpload> => {
|
||||
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<PreparedUpload> => {
|
||||
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<Response> => {
|
||||
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<Response> => {
|
||||
* 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<Response> => {
|
||||
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<Response> => {
|
||||
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<Response> => {
|
||||
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 });
|
||||
|
||||
@@ -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<string, string> = {
|
||||
'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 },
|
||||
);
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user