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) <noreply@anthropic.com>
This commit is contained in:
Claude
2026-07-28 17:59:49 +07:00
parent 5e29589f1a
commit e2de8c4245
3 changed files with 644 additions and 0 deletions
+108
View File
@@ -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<LoginResponse> => {
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<LogoutResponse> => {
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<UserInfoResponse | null> => {
if (!isAuthEnabled(deps.config.adminApiToken)) {
return null;
}
if (!session) {
return null;
}
return {
username: session.username,
expiresAt: session.expiresAt?.toISOString() ?? null,
};
};
}
+176
View File
@@ -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<GetFileUseCaseDeps, 'fileRepo'>) {
return async (publicId: string): Promise<GetFileInfoResult> => {
const file = await deps.fileRepo.findByPublicId(publicId);
if (!file) {
return { found: false };
}
return {
found: true,
publicId: file.publicId,
fileName: file.fileName,
mimeType: file.mimeType,
sizeBytes: file.sizeBytes,
fileType: file.fileType,
createdAt: formatCreatedAtForInfo(file.createdAt),
};
};
}
/**
* Formats a date-like value into an ISO-8601 string.
*
* @param date - A Date instance, date string, or numeric timestamp.
* @returns The ISO-8601 string.
*/
const formatCreatedAtForInfo = (date: Date | string | number): string => {
if (date instanceof Date) return date.toISOString();
return new Date(date).toISOString();
};
/**
* Creates a factory function for the get-file retrieval use case.
*
* Determines how a file should be delivered to the client:
* - **redirect**: For regular (non-chunked, non-archive) files — returns a
* Telegram CDN redirect URL.
* - **chunked**: For files stored across multiple Telegram parts — returns
* the file entity so the caller can build a multi-part streaming response.
* - **archive-entry**: For files stored inside a Telegram archive (zip) —
* returns the archive's Telegram metadata and the entry name so the caller
* can extract and stream the entry.
*
* @param deps - The injected dependencies.
* @returns An async function accepting a public ID and returning a retrieval result.
*/
export function createGetFileUseCase(deps: GetFileUseCaseDeps) {
return async (publicId: string): Promise<FileRetrievalResult | null> => {
const file = await deps.fileRepo.findByPublicId(publicId);
if (!file) {
return null;
}
// Chunked file — return the entity for multi-part response building
if (file.storageBackend === 'chunked') {
return { type: 'chunked', file };
}
// Archive entry — resolve the archive's Telegram location
const archiveEntryName = file.archiveEntryName;
if (archiveEntryName) {
const archiveFileId = file.archiveTelegramFileId || file.telegramFileId;
const archiveInfo = await deps.telegramService.getFileInfo(archiveFileId);
return { type: 'archive-entry', file, archiveInfo, entryName: archiveEntryName };
}
// Regular file — resolve Telegram CDN path 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 };
};
}
+360
View File
@@ -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<Buffer> => {
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<ChunkedUploadResult> => {
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<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}`,
};
}
// 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}`,
};
};
}