14 Commits
Author SHA1 Message Date
semantic-release-bot 083fbe1a6f chore(release): 1.2.3
## [1.2.3](https://github.com/asepharyana/TeleUploader/compare/v1.2.2...v1.2.3) (2026-09-18)
2026-09-18 17:55:15 +00:00
asepharyana b438981231 ci(release): trigger patch release for refactor/test/build/ci/docs commits 2026-09-19 00:54:30 +07:00
asepharyana 890444fed6 chore(plan): fix stale src/db doc refs + register s3-routing test
- 8 repository/port doc comments pointed at long-gone src/db/*
  layout; now reference the real drizzle repository paths
- test:unit and test:s3 now include test/s3-routing.test.ts
  (was passing but never wired into any script)
2026-09-14 18:19:22 +07:00
asepharyana 3a9052d551 refactor(plan): delete unused DTO modules bucket, file, s3
Zero importers in src/ and test/ (verified by grep). Live DTOs:
auth (authenticate use-case + controller) and upload (upload-file
use-case). Build + biome green; full test:unit green.
2026-09-14 18:17:55 +07:00
asepharyana e610f1ca21 refactor(plan): unify bucket validation + lock root-verb routing
- web-api create-bucket and s3 virtual-host now use BucketNameSchema
  (single canonical validator; no inline regex left in src/)
- routes '/': HEAD/DELETE/POST now check shouldHandleS3 with headers
  like every other branch (non-S3 -> 404, never S3-direct blind)
- s3-routing.test.ts: 3 new tests locking the fixed root verbs
2026-09-14 18:17:18 +07:00
asepharyana da7fca9fdc refactor(plan): delete dead use-cases s3-object, manage-bucket, multipart-upload, get-file
Verified zero importers in src/ and test/ (only a doc comment
mentions manage-bucket). Live use-cases: authenticate (auth
controller) and upload-file (upload + web-api controllers); S3
handlers call repositories directly. Build + biome + full
test:unit (28 files) green.
2026-09-14 18:12:21 +07:00
asepharyana 2d20365d56 test(integration): fix stale assertions exposed by refactor
- files.test.ts: 302-redirect expectation replaced with 200-proxy
  assertion (bot token only in server-side fetch, never Location)
- deploy-config.test.ts: retarget from deleted .gitea workflow to
  .github/workflows/deploy.yml (CI migrated to GitHub Actions in
  811a688); asserts nix build + VPS deploy shape
2026-09-14 18:07:23 +07:00
asepharyana bd4cf8b285 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/)
2026-09-14 18:03:37 +07:00
asepharyana 7d5b8154a4 Merge branch 'refactor/fase1c-fondasi' into refactor/teleuploader-full 2026-09-14 17:43:19 +07:00
asepharyana cba2160acb test(fase1-s3): retarget source-reading asserts to split s3/ modules
PUT streaming assert now reads s3/s3-object-write.ts and UploadPart
assert reads s3/s3-multipart-handlers.ts (s3-controller.ts is a thin
facade since the split). No behavior change.
2026-09-14 17:28:22 +07:00
asepharyana b164b8826f refactor(fase1c): routing S3 fixes + auth/swagger/index/env
- routes: GET / teruskan headers ke shouldHandleS3; OPTIONS jawab
  204 CORS generik bila bukan S3, handleS3Direct bila S3; komentar
  bypass rate-limit S3 dipertahankan (registry abort pada 429)
- auth-controller: readLoginBody via LoginBodySchema; handleMe sederhanakan
  (getAuthSession sudah cek bearer); import AuthSession dari dto
- swagger: pindah src/routes -> src/interfaces/http/swagger; tambah path
  auth, /api/v1, /{bucket}, /{bucket}/{key}; version dari config.appVersion
- index: unref ketiga setInterval agar tak menahan process
- env: PORT fail-fast via PositiveIntSchema (default 4000 bila tak diset);
  log config turun ke debug
- metrics: dokumentasikan uploadThroughput/queueSize/botUtilization
- test baru test/s3-routing.test.ts (GET / S3 vs home, OPTIONS 204 CORS)
2026-09-14 17:26:03 +07:00
asepharyana d323260ed5 refactor(fase1-s3): split s3-controller into s3/ modules + guards 2026-09-14 17:25:07 +07:00
asepharyana 5d01a9405f refactor(fase0): shared scaffold — crypto, file-url, schemas, test splits
- Add zod dep; new src/shared/validation/schemas.ts (BucketName,
  JsonUploadPayload, LoginBody, DeleteObjects, CompleteMultipart,
  clampMaxKeys, parseOrNull) + test/validation.test.ts
- Dedup timingSafeCompare -> src/shared/utils/crypto.ts (auth
  middleware, authenticate use-case, s3/auth now import it)
- Dedup AuthSession -> single type in dto/auth.ts
- Dedup telegram file URL builder + filename sanitizer to shared
  modules (file-controller now imports the canonical ones)
- Remove duplicate controllers/home.html (canonical: src/home.html)
- Fix stale bootstrap mocks to real module paths; bootstrap now
  asserts public GET vs guarded POST separately
- Fix health assertion to include version field
- package.json: test -> test:unit alias; new test:quarantine for
  network/live tests; register previously unlisted test files
2026-09-14 17:09:52 +07:00
asepharyana f4b30cef0c chore(gitignore): ignore Ruflo runtime, secrets, and scaffolding 2026-09-14 16:13:46 +07:00
63 changed files with 3415 additions and 3796 deletions
+168 -1
View File
@@ -6,7 +6,174 @@
"Bun(bunx *)",
"Bun(bun run *)",
"Bun(bun build *)",
"Bun(bun install)"
"Bun(bun install)",
"Bash(npx @claude-flow*)",
"Bash(npx claude-flow*)",
"Bash(node .claude/*)",
"mcp__claude-flow__*"
]
},
"hooks": {
"PreToolUse": [
{
"matcher": "Bash",
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" pre-bash'",
"timeout": 5000
}
]
},
{
"matcher": "Write|Edit|MultiEdit",
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" pre-edit'",
"timeout": 5000
}
]
}
],
"PostToolUse": [
{
"matcher": "Write|Edit|MultiEdit",
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" post-edit'",
"timeout": 10000
}
]
},
{
"matcher": "Bash",
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" post-bash'",
"timeout": 5000
}
]
}
],
"UserPromptSubmit": [
{
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" route'",
"timeout": 10000
}
]
}
],
"SessionStart": [
{
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" session-restore'",
"timeout": 15000
},
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/auto-memory-hook.mjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/auto-memory-hook.mjs\" import'",
"timeout": 8000
}
]
}
],
"SessionEnd": [
{
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" session-end'",
"timeout": 10000
}
]
}
],
"Stop": [
{
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/auto-memory-hook.mjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/auto-memory-hook.mjs\" sync'",
"timeout": 10000
}
]
}
],
"PreCompact": [
{
"matcher": "manual",
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" compact-manual'"
},
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" session-end'",
"timeout": 5000
}
]
},
{
"matcher": "auto",
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" compact-auto'"
},
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" session-end'",
"timeout": 6000
}
]
}
],
"SubagentStart": [
{
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" status'",
"timeout": 3000
}
]
}
],
"SubagentStop": [
{
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" post-task'",
"timeout": 5000
}
]
}
],
"Notification": [
{
"hooks": [
{
"type": "command",
"command": "sh -c 'D=\"${CLAUDE_PROJECT_DIR:-.}\"; [ -f \"$D/.claude/helpers/hook-handler.cjs\" ] || D=\"${HOME}\"; exec node \"$D/.claude/helpers/hook-handler.cjs\" notify'",
"timeout": 3000
}
]
}
]
},
"env": {
"CLAUDE_CODE_EXPERIMENTAL_AGENT_TEAMS": "1",
"CLAUDE_FLOW_V3_ENABLED": "true",
"CLAUDE_FLOW_HOOKS_ENABLED": "true"
}
}
+33
View File
@@ -1,5 +1,6 @@
# dependencies (bun install)
node_modules
package-lock.json
# output
out
@@ -37,3 +38,35 @@ S3_GUIDE.md
# Finder (MacOS) folder config
.DS_Store
result
# Ruflo / Claude Flow local state (scaffolding kept on disk, never committed)
# Secrets & configs with credentials
claude-flow.config.json
claude-flow.config.json.bak
.env.*.local
# Runtime data & databases
*.db
*.db-shm
*.db-wal
ruvector.db
data/
.swarm/
.claude-flow/
.claude/*.db
.claude/.proven-config-version
.claude/proven-config.json
# Generated scaffolding (re-creatable via ruflo init; force-add if you want it tracked)
.agents/
AGENTS.md
.claude/agents/
.claude/commands/
.claude/helpers/
.claude/skills/
# Codex local configuration
.codex/
# Planning files (per-repo)
.planning/
+12 -1
View File
@@ -1,7 +1,18 @@
{
"branches": ["main"],
"plugins": [
"@semantic-release/commit-analyzer",
[
"@semantic-release/commit-analyzer",
{
"releaseRules": [
{ "type": "refactor", "release": "patch" },
{ "type": "test", "release": "patch" },
{ "type": "build", "release": "patch" },
{ "type": "ci", "release": "patch" },
{ "type": "docs", "release": "patch" }
]
}
],
"@semantic-release/release-notes-generator",
[
"@semantic-release/changelog",
+2
View File
@@ -1,3 +1,5 @@
## [1.2.3](https://github.com/asepharyana/TeleUploader/compare/v1.2.2...v1.2.3) (2026-09-18)
## [1.2.2](https://github.com/asepharyana/TeleUploader/compare/v1.2.1...v1.2.2) (2026-09-05)
+3
View File
File diff suppressed because one or more lines are too long
+1 -1
View File
@@ -27,7 +27,7 @@
# TeleUploader package
teleuploader = pkgs.stdenvNoCC.mkDerivation rec {
pname = "teleuploader";
version = "1.2.2";
version = "1.2.3";
src = ./.;
+7 -4
View File
File diff suppressed because one or more lines are too long
-47
View File
@@ -1,47 +0,0 @@
/**
* Input for creating a new bucket.
*/
export interface CreateBucketInput {
/** Bucket name (must match S3 naming rules: 3-63 chars, lowercase, no underscore) */
name: string;
}
/**
* Single bucket representation returned by bucket endpoints.
*/
export interface BucketResponse {
/** Bucket UUID */
id: string;
/** Bucket name */
name: string;
/** ISO-8601 timestamp of when the bucket was created */
createdAt: string;
/** Number of non-deleted objects in the bucket */
objectCount?: number;
}
/**
* Response payload for the list-buckets endpoint.
*/
export interface BucketListResponse {
/** Array of buckets */
buckets: BucketResponse[];
}
/**
* Response payload for bucket creation.
*/
export interface CreateBucketResponse {
/** Bucket UUID */
id: string;
/** Bucket name */
name: string;
}
/**
* Response payload for bucket deletion.
*/
export interface DeleteBucketResponse {
/** Whether the deletion succeeded */
success: boolean;
}
-68
View File
@@ -1,68 +0,0 @@
/**
* Public file information response returned by the file-info endpoint.
* Mirrors the JSON shape of GET /file/:publicId/info.
*/
export interface FileInfoResponse {
/** Public, shareable identifier (nanoid) */
public_id: string;
/** Stored file name */
file_name: string;
/** MIME type of the stored file */
mime_type: string;
/** File size in bytes */
size_bytes: number;
/** High-level file category (e.g. "document", "photo") */
file_type: string;
/** ISO-8601 timestamp of when the file record was created */
created_at: string;
}
/**
* Summary-level file metadata used internally for constructing
* upload responses and object listing entries.
*/
export interface FileMetadata {
/** Public, shareable identifier (nanoid) */
publicId: string;
/** Telegram file identifier used to retrieve the file from Telegram CDN */
telegramFileId: string;
/** Telegram unique file identifier (persists across re‑uploads) */
telegramFileUniqueId: string;
/** Chat ID where the file or archive was stored */
storageChatId: number;
/** Message ID of the stored file or archive */
storageMessageId: number;
/** Stored file name */
fileName: string;
/** MIME type of the stored file */
mimeType: string;
/** File size in bytes */
sizeBytes: number;
/** High-level file category */
fileType: string;
/** Telegram user ID of the uploader; 0 when unknown or system */
uploaderId: number;
/** Timestamp of file record creation */
createdAt: Date | string | number;
}
/**
* Upload response shape returned to API callers.
* Mirrors the JSON output of the /api/upload endpoint.
*/
export interface UploadResponse {
/** Public, shareable identifier */
public_id: string;
/** Stored file name */
file_name: string;
/** MIME type */
mime_type: string;
/** File size in bytes */
size_bytes: number;
/** High-level file category */
file_type: string;
/** ISO-8601 creation timestamp */
created_at: string;
/** Public download URL */
download_url: string;
}
-87
View File
@@ -1,87 +0,0 @@
/**
* A single S3 object as it appears in listing results.
*/
export interface S3ObjectResponse {
/** The object key (full path within the bucket) */
key: string;
/** Stored file name (basename of the key) */
fileName: string;
/** MIME type of the stored object */
mimeType: string;
/** Object size in bytes */
sizeBytes: number;
/** High-level file category */
fileType: string;
/** SHA-256 hex digest of the object content */
etag: string | null;
/** ISO-8601 timestamp of last modification */
lastModified: string;
/** Public download URL */
downloadUrl: string;
}
/**
* Response payload for S3 ListObjectsV1 / ListObjectsV2.
*/
export interface S3ListObjectsResponse {
/** Array of object summaries */
objects: S3ObjectResponse[];
/** Common prefixes when a delimiter was used (e.g. "folder/" entries) */
prefixes: string[];
/** Whether more results are available */
isTruncated: boolean;
/** Token to pass as continuation-token to retrieve the next page */
nextContinuationToken: string | null;
}
/**
* Input for the copy-object operation (Web API v1).
*/
export interface S3CopyObjectInput {
/** Source object key within the same or source bucket */
sourceKey: string;
/** Destination bucket name; defaults to the source bucket when omitted */
destBucket?: string;
/** Destination object key */
destKey: string;
}
/**
* Response payload for the copy-object operation.
*/
export interface S3CopyObjectResponse {
/** Source object key that was copied */
sourceKey: string;
/** Destination object key */
destKey: string;
/** Destination bucket name */
destBucket: string;
}
/**
* Summary of a multipart upload in listing results.
*/
export interface S3MultipartUploadResponse {
/** The object key being uploaded */
key: string;
/** Upload identifier (nanoid) */
uploadId: string;
/** ISO-8601 timestamp when the upload was initiated */
initiatedAt: Date;
/** Identifier string of the upload initiator */
initiatedBy: string;
}
/**
* Summary of a single part within a multipart upload.
*/
export interface S3MultipartPartResponse {
/** 1-indexed part number */
partNumber: number;
/** ETag of the part content */
etag: string;
/** Part size in bytes */
sizeBytes: number;
/** ISO-8601 timestamp when the part was stored */
createdAt: Date;
}
+24 -2
View File
@@ -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 -19
View File
@@ -1,4 +1,4 @@
import { timingSafeEqual } from 'node:crypto';
import { timingSafeCompare } from '../../shared/utils/crypto';
import type {
AuthSession,
LoginInput,
@@ -23,24 +23,6 @@ export interface AuthenticateUseCaseDeps {
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.
*
-176
View File
@@ -1,176 +0,0 @@
import type { File } from '../../domain/entities/file';
import type { IFileRepository } from '../../domain/ports/file-repository';
import type { ITelegramService, TelegramFileInfo } from '../../domain/ports/telegram-service';
/**
* 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 };
};
}
-156
View File
@@ -1,156 +0,0 @@
import type { Bucket } from '../../domain/entities/bucket';
import type { IBucketRepository } from '../../domain/ports/bucket-repository';
import type { IFileRepository } from '../../domain/ports/file-repository';
/**
* S3 bucket name validation regex.
*
* Bucket names must be 3-63 characters, start/end with a lowercase letter or
* digit, and contain only lowercase letters, digits, dots, and hyphens.
*/
const BUCKET_NAME_REGEX = /^[a-z0-9][a-z0-9.-]{1,61}[a-z0-9]$/;
/**
* Error type for bucket-level application errors that carry an S3-compatible
* error code and an HTTP status suggestion.
*/
export class BucketError extends Error {
/** S3-compatible error code (e.g. "NoSuchBucket", "BucketAlreadyExists"). */
readonly code: string;
/** Suggested HTTP status code. */
readonly status: number;
/**
* @param code - The S3 error code.
* @param message - Human-readable error description.
* @param status - Suggested HTTP status.
*/
constructor(code: string, message: string, status: number) {
super(message);
this.name = 'BucketError';
this.code = code;
this.status = status;
}
}
/**
* Describes a bucket entry returned by the list-buckets use case,
* enriched with the current count of non-deleted objects.
*/
export interface BucketWithCount {
/** The bucket domain entity. */
bucket: Bucket;
/** Number of non-deleted objects in the bucket. */
objectCount: number;
}
/** Dependencies required by the manage-bucket use case factories. */
export interface ManageBucketDeps {
/** Bucket repository for CRUD operations. */
bucketRepo: IBucketRepository;
/** File repository for counting and checking objects within buckets. */
fileRepo: IFileRepository;
}
/**
* Creates a use case that lists all buckets together with their object counts.
*
* @param deps - The injected dependencies.
* @returns An async function that returns a list of buckets with counts.
*/
export function createListBucketsUseCase(deps: ManageBucketDeps) {
return async (): Promise<BucketWithCount[]> => {
const buckets = await deps.bucketRepo.list();
const results = await Promise.all(
buckets.map(async (bucket) => ({
bucket,
objectCount: await deps.fileRepo.countByBucket(bucket.id),
})),
);
return results;
};
}
/**
* Creates a use case that retrieves a single bucket by its name.
*
* @param deps - The injected dependencies.
* @returns An async function accepting a bucket name and returning the
* bucket, or `null` when not found.
*/
export function createGetBucketUseCase(deps: ManageBucketDeps) {
return async (name: string): Promise<Bucket | null> => {
return deps.bucketRepo.findByName(name);
};
}
/**
* Creates a use case that creates a new bucket.
*
* Validates the bucket name format and checks for duplicates before
* persisting.
*
* @param deps - The injected dependencies.
* @returns An async function accepting a bucket name and returning the
* newly created bucket.
* @throws {BucketError} When the name is invalid or the bucket already exists.
*/
export function createCreateBucketUseCase(deps: ManageBucketDeps) {
return async (name: string): Promise<Bucket> => {
if (!BUCKET_NAME_REGEX.test(name)) {
throw new BucketError('InvalidBucketName', 'The specified bucket is not valid.', 400);
}
const existing = await deps.bucketRepo.findByName(name);
if (existing) {
throw new BucketError(
'BucketAlreadyExists',
'The requested bucket name is not available.',
409,
);
}
return deps.bucketRepo.create(name);
};
}
/**
* Creates a use case that deletes a bucket.
*
* Ensures the bucket exists and is empty (no non-deleted objects) before
* proceeding with deletion.
*
* @param deps - The injected dependencies.
* @returns An async function accepting a bucket name. Returns `true` when
* the bucket was deleted, throws when the bucket is missing or
* not empty.
* @throws {BucketError} When the bucket does not exist or is not empty.
*/
export function createDeleteBucketUseCase(deps: ManageBucketDeps) {
return async (name: string): Promise<boolean> => {
const bucket = await deps.bucketRepo.findByName(name);
if (!bucket) {
throw new BucketError('NoSuchBucket', 'The specified bucket does not exist.', 404);
}
const objectCount = await deps.fileRepo.countByBucket(bucket.id);
if (objectCount > 0) {
throw new BucketError('BucketNotEmpty', 'The bucket you tried to delete is not empty.', 409);
}
return deps.bucketRepo.delete(name);
};
}
/**
* Creates a use case that checks whether a bucket exists.
*
* @param deps - The injected dependencies.
* @returns An async function accepting a bucket name and returning `true`
* when the bucket exists.
*/
export function createBucketExistsUseCase(deps: ManageBucketDeps) {
return async (name: string): Promise<boolean> => {
return deps.bucketRepo.exists(name);
};
}
@@ -1,368 +0,0 @@
import { nanoid } from 'nanoid';
import { buildNewFile } from '../../domain/entities/file-factory';
import type { MultipartUpload } from '../../domain/entities/multipart';
import type { IBucketRepository } from '../../domain/ports/bucket-repository';
import type { IFileRepository } from '../../domain/ports/file-repository';
import type { IMultipartRepository } from '../../domain/ports/multipart-repository';
import type { ITelegramService } from '../../domain/ports/telegram-service';
import { computeHash, DEFAULT_FILE_TYPE } from '../../shared/utils/file';
// ─── Types ──────────────────────────────────────────────────────────
/**
* A single part reference as submitted in a complete-multipart-upload request.
*/
export interface CompletePartInput {
/** 1-based part number. */
partNumber: number;
/** ETag returned when the part was uploaded. */
etag: string;
}
/**
* Result of initiating a multipart upload.
*/
export interface InitiateMultipartResult {
/** The generated upload identifier (nanoid). */
uploadId: string;
/** The bucket name. */
bucket: string;
/** The S3 object key. */
key: string;
}
/**
* Result of uploading a single part.
*/
export interface UploadPartResult {
/** ETag of the uploaded part (SHA-256 hex digest). */
etag: string;
}
/**
* Result of completing a multipart upload.
*/
export interface CompleteMultipartResult {
/** Public-facing unique identifier of the created file record. */
publicId: string;
/** The S3 location URL of the completed object. */
location: string;
/** Combined ETag (all part etags joined by hyphens). */
etag: string;
/** Total object size in bytes. */
sizeBytes: number;
}
/**
* Summary of a single part within a multipart upload used in listing results.
*/
export interface PartSummary {
/** 1-based part number. */
partNumber: number;
/** ETag of the part content. */
etag: string;
/** Part size in bytes. */
sizeBytes: number;
/** ISO-8601 timestamp when the part was stored. */
createdAt: Date;
}
/**
* Result of listing multipart uploads within a bucket.
*/
export interface ListMultipartUploadsResult {
/** Array of in-progress upload summaries. */
uploads: MultipartUpload[];
/** Whether more results are available. */
isTruncated: boolean;
/** Marker for the next page, or null when not truncated. */
nextKeyMarker: string | null;
}
// ─── Config ─────────────────────────────────────────────────────────
/** Subset of application configuration consumed by the multipart use cases. */
export interface MultipartConfig {
/** Maximum chunk size in bytes for Telegram uploads (part size limit). */
telegramChunkSizeBytes: number;
/** Telegram chat ID where part data is stored. */
storageChatId: number;
/** Server base URL for constructing location URLs. */
baseUrl: string;
}
/** Dependencies required by the multipart upload use case factories. */
export interface MultipartDeps {
/** Bucket repository for bucket lookups. */
bucketRepo: IBucketRepository;
/** File repository for creating the final file record on completion. */
fileRepo: IFileRepository;
/** Multipart repository for managing upload sessions and parts. */
multipartRepo: IMultipartRepository;
/** Telegram service for forwarding part data to storage. */
telegramService: ITelegramService;
/** Application configuration subset. */
config: MultipartConfig;
}
// ─── Use Case Factories ─────────────────────────────────────────────
/**
* Creates a use case that initiates an S3 multipart upload.
*
* Validates the bucket exists and creates a new multipart upload session.
*
* @param deps - The injected dependencies.
* @returns An async function accepting bucket name and object key, returning
* the upload initiation result, or `null` when the bucket is not found.
*/
export function createInitiateMultipartUploadUseCase(deps: MultipartDeps) {
return async (bucketName: string, key: string): Promise<InitiateMultipartResult | null> => {
const bucket = await deps.bucketRepo.findByName(bucketName);
if (!bucket) return null;
const uploadId = await deps.multipartRepo.create(bucket.id, key, 's3');
return { uploadId, bucket: bucketName, key };
};
}
/**
* Creates a use case that uploads a single part of a multipart upload.
*
* Validates the part number range (1-10000), checks the upload session exists
* and matches the expected key, checks part size against the configured limit,
* forwards the part data to Telegram storage, and persists the part record.
*
* @param deps - The injected dependencies.
* @returns An async function accepting upload details and part data, returning
* the part ETag, or `null` when the upload session is not found.
*/
export function createUploadPartUseCase(deps: MultipartDeps) {
return async (input: {
/** Bucket name for the multipart upload. */
bucketName: string;
/** S3 object key for the multipart upload. */
key: string;
/** Upload identifier returned by initiate. */
uploadId: string;
/** 1-based part number (1-10000). */
partNumber: number;
/** Raw part data. */
body: Buffer;
}): Promise<UploadPartResult | null> => {
if (input.partNumber < 1 || input.partNumber > 10000) {
throw new MultipartError(
'InvalidArgument',
'Part number must be an integer between 1 and 10000',
400,
);
}
const multipart = await deps.multipartRepo.findById(input.uploadId);
if (!multipart || multipart.s3Key !== input.key) {
return null;
}
if (input.body.byteLength > deps.config.telegramChunkSizeBytes) {
throw new MultipartError(
'EntityTooLarge',
`Your proposed upload part size (${input.body.byteLength} bytes) exceeds the maximum allowed part size (${deps.config.telegramChunkSizeBytes} bytes) for this storage backend. Use smaller part sizes.`,
400,
);
}
const forwardResult = await deps.telegramService.forwardToStorage(
input.body,
`mp-${input.uploadId}-part-${input.partNumber}`,
'document',
);
const etag = computeHash(input.body);
await deps.multipartRepo.insertPart({
uploadId: input.uploadId,
partNumber: input.partNumber,
telegramFileId: forwardResult.telegramFileId,
telegramFileUniqueId: forwardResult.telegramFileUniqueId,
storageMessageId: forwardResult.storageMessageId,
sizeBytes: input.body.byteLength,
etag,
});
return { etag: `"${etag}"` };
};
}
/**
* Error type for multipart-level application errors.
*/
export class MultipartError extends Error {
/** S3-compatible error code. */
readonly code: string;
/** Suggested HTTP status code. */
readonly status: number;
/**
* @param code - The S3 error code.
* @param message - Human-readable error description.
* @param status - Suggested HTTP status.
*/
constructor(code: string, message: string, status: number) {
super(message);
this.name = 'MultipartError';
this.code = code;
this.status = status;
}
}
/**
* Creates a use case that completes an S3 multipart upload.
*
* Validates the submitted part list (all parts must be present and in ascending
* order), creates the final file record referencing the first part's Telegram
* data, marks the upload session as completed, and returns the combined result.
*
* @param deps - The injected dependencies.
* @returns An async function accepting upload details and submitted parts,
* returning the completion result, or `null` when the upload session
* is not found.
*/
export function createCompleteMultipartUploadUseCase(deps: MultipartDeps) {
return async (input: {
/** Bucket name for the multipart upload. */
bucketName: string;
/** S3 object key for the multipart upload. */
key: string;
/** Upload identifier. */
uploadId: string;
/** Parts submitted by the client (in ascending part number order). */
parts: CompletePartInput[];
}): Promise<CompleteMultipartResult | null> => {
const multipart = await deps.multipartRepo.findById(input.uploadId);
if (!multipart) return null;
const storedParts = await deps.multipartRepo.listParts(input.uploadId);
// Validate ascending part order
const partNumbers = input.parts.map((p) => p.partNumber);
if (partNumbers.length > 1 && partNumbers.some((n, i) => i > 0 && n <= partNumbers[i - 1])) {
throw new MultipartError(
'InvalidPartOrder',
'The list of parts was not in ascending order.',
400,
);
}
// Validate part count matches
if (input.parts.length !== storedParts.length) {
throw new MultipartError(
'InvalidPart',
'One or more specified parts could not be found.',
400,
);
}
const totalSize = storedParts.reduce((sum, p) => sum + p.sizeBytes, 0);
const firstPart = storedParts[0];
if (!firstPart) {
throw new MultipartError('InternalError', 'Multipart object has no parts.', 500);
}
const publicId = nanoid();
await deps.fileRepo.create(
buildNewFile({
publicId,
telegramFileId: firstPart.telegramFileId,
telegramFileUniqueId: firstPart.telegramFileUniqueId,
storageChatId: deps.config.storageChatId,
storageMessageId: firstPart.storageMessageId,
fileName: input.key.split('/').pop() || 'file',
mimeType: 'application/octet-stream',
sizeBytes: totalSize,
fileType: DEFAULT_FILE_TYPE,
storageBackend: 'telegram',
bucketId: multipart.bucketId,
s3Key: input.key,
multipartUploadId: input.uploadId,
}),
);
await deps.multipartRepo.complete(input.uploadId);
const location = `${deps.config.baseUrl}/${input.bucketName}/${input.key}`;
const combinedEtag = storedParts.map((p) => p.etag).join('-');
return {
publicId,
location,
etag: combinedEtag,
sizeBytes: totalSize,
};
};
}
/**
* Creates a use case that aborts an S3 multipart upload.
*
* @param deps - The injected dependencies.
* @returns An async function accepting an upload identifier, returning `true`
* when the upload was aborted, or `null` when the upload session
* is not found.
*/
export function createAbortMultipartUploadUseCase(deps: MultipartDeps) {
return async (uploadId: string): Promise<boolean | null> => {
const multipart = await deps.multipartRepo.findById(uploadId);
if (!multipart) return null;
await deps.multipartRepo.abort(uploadId);
return true;
};
}
/**
* Creates a use case that lists in-progress multipart uploads within a bucket.
*
* @param deps - The injected dependencies.
* @returns An async function accepting query parameters and returning the
* listing result, or `null` when the bucket is not found.
*/
export function createListMultipartUploadsUseCase(deps: MultipartDeps) {
return async (input: {
/** Bucket name to list uploads from. */
bucketName: string;
/** Maximum number of uploads to return (clamped 1-1000). */
maxUploads: number;
/** Return only uploads whose S3 key is strictly greater than this, or null. */
keyMarker: string | null;
}): Promise<ListMultipartUploadsResult | null> => {
const bucket = await deps.bucketRepo.findByName(input.bucketName);
if (!bucket) return null;
return deps.multipartRepo.listByBucket(bucket.id, input.maxUploads, input.keyMarker);
};
}
/**
* Creates a use case that lists parts of a specific multipart upload.
*
* @param deps - The injected dependencies.
* @returns An async function accepting an upload identifier and returning the
* list of parts, or `null` when the upload session is not found.
*/
export function createListPartsUseCase(deps: MultipartDeps) {
return async (uploadId: string): Promise<PartSummary[] | null> => {
const multipart = await deps.multipartRepo.findById(uploadId);
if (!multipart) return null;
const parts = await deps.multipartRepo.listParts(uploadId);
return parts.map((p) => ({
partNumber: p.partNumber,
etag: p.etag,
sizeBytes: p.sizeBytes,
createdAt: p.createdAt,
}));
};
}
-165
View File
@@ -1,165 +0,0 @@
import type { File } from '../../domain/entities/file';
import type { IBucketRepository } from '../../domain/ports/bucket-repository';
import type { IFilePartRepository } from '../../domain/ports/file-part-repository';
import type { IFileRepository } from '../../domain/ports/file-repository';
import type { IMultipartRepository } from '../../domain/ports/multipart-repository';
import type { ITelegramService, TelegramFileInfo } from '../../domain/ports/telegram-service';
import type { CompressionAlgorithm } from '../../shared/utils/compress';
// ─── Types ──────────────────────────────────────────────────────────
/**
* A single part source for building a multi-part streaming response.
* Each part corresponds to a Telegram-stored file chunk.
*/
export interface ObjectPartSource {
/** Telegram file identifier for retrieving this part. */
telegramFileId: string;
/** Telegram CDN URL for downloading this part. */
telegramUrl: string;
/** Original size of this part in bytes. */
sizeBytes: number;
/** 1-based part number within the object. */
partNumber: number;
/** Stored (post-compression) size in bytes, when applicable. */
storedSizeBytes?: number;
/** Compression algorithm applied, or null if uncompressed. */
compressionAlgorithm?: CompressionAlgorithm;
}
/**
* A regular (direct) S3 object resolved to a Telegram CDN URL.
*/
export interface DirectObjectResult {
/** Discriminant. */
type: 'direct';
/** The resolved file entity. */
file: File;
/** Full Telegram CDN URL for downloading the object. */
telegramUrl: string;
/** Telegram file metadata. */
fileInfo: TelegramFileInfo;
}
/**
* A chunked S3 object stored across multiple Telegram file parts.
*/
export interface ChunkedObjectResult {
/** Discriminant. */
type: 'chunked';
/** The resolved file entity. */
file: File;
}
/**
* An S3 object assembled from a completed multipart upload.
*/
export interface MultipartObjectResult {
/** Discriminant. */
type: 'multipart';
/** The resolved file entity. */
file: File;
/** Resolved part sources with Telegram CDN URLs. */
parts: ObjectPartSource[];
}
/**
* Discriminated union of all possible S3 get-object outcomes.
*/
export type GetObjectResult = DirectObjectResult | ChunkedObjectResult | MultipartObjectResult;
/**
* Result of an S3 put-object operation.
*/
export interface PutObjectResult {
/** SHA-256 hex digest of the object content. */
etag: string;
}
/**
* Result of an S3 copy-object operation.
*/
export interface CopyObjectResult {
/** SHA-256 hex digest of the source object content. */
etag: string;
/** ISO-8601 timestamp of the copy operation. */
lastModified: string;
}
/**
* A single S3 object as returned in listing results.
*/
export interface ListObjectEntry {
/** The object key (full path within the bucket). */
key: string;
/** Object size in bytes. */
sizeBytes: number;
/** SHA-256 hex digest or fallback identifier. */
etag: string;
/** ISO-8601 timestamp of last modification. */
lastModified: string;
/** MIME type of the stored object. */
mimeType: string;
}
/**
* Result of an S3 list-objects operation (both V1 and V2).
*/
export interface ListObjectsResult {
/** Array of object summaries. */
objects: ListObjectEntry[];
/** Common prefixes when a delimiter was used. */
prefixes: string[];
/** Whether more results are available. */
isTruncated: boolean;
/** The last key in the returned page, for use as the next marker. */
nextMarker: string | null;
}
/**
* Head-object metadata.
*/
export interface HeadObjectMetadata {
/** MIME type of the object. */
contentType: string;
/** Object size in bytes. */
contentLength: number;
/** Entity tag (SHA-256 hex digest or fallback). */
etag: string;
/** ISO-8601 timestamp of last modification. */
lastModified: string;
}
// ─── Config ─────────────────────────────────────────────────────────
/** Subset of application configuration consumed by the s3-object use cases. */
export interface S3ObjectConfig {
/** Maximum chunk size in bytes for Telegram chunked uploads. */
telegramChunkSizeBytes: number;
/** Whether gzip compression is enabled for chunked uploads. */
compressChunkedUploads: boolean;
/** Minimum chunk size in bytes below which compression is skipped. */
chunkCompressionMinSizeBytes: number;
/** Telegram chat ID where file parts are stored. */
storageChatId: number;
/** Server base URL for constructing download links. */
baseUrl: string;
/** Whether to proxy S3 GET requests through the server. */
proxyS3Get: boolean;
}
/** Dependencies required by the s3-object use case factories. */
export interface S3ObjectDeps {
/** Bucket repository for bucket lookups. */
bucketRepo: IBucketRepository;
/** File repository for object CRUD operations. */
fileRepo: IFileRepository;
/** File-part repository for chunked upload part records. */
filePartRepo: IFilePartRepository;
/** Multipart repository for resolving multipart-upload objects. */
multipartRepo: IMultipartRepository;
/** Telegram service for uploading and resolving file metadata. */
telegramService: ITelegramService;
/** Application configuration subset. */
config: S3ObjectConfig;
}
+102 -80
View File
@@ -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}`,
};
};
}
+1 -1
View File
@@ -3,7 +3,7 @@ import type { Bucket } from '../entities/bucket';
/**
* Repository interface for Bucket entity persistence.
*
* Abstracts the bucket CRUD operations currently in `src/db/buckets.ts`.
* Abstracts the bucket CRUD operations currently in `infrastructure/persistence/repositories/bucket-repository.ts`.
*/
export interface IBucketRepository {
/**
+1 -1
View File
@@ -3,7 +3,7 @@ import type { FilePart, NewFilePart } from '../entities/file-part';
/**
* Repository interface for FilePart entity persistence.
*
* Abstracts the file-part operations currently in `src/db/file-parts.ts`.
* Abstracts the file-part operations currently in `infrastructure/persistence/repositories/file-part-repository.ts`.
* File parts represent the chunks of a large file stored across multiple
* Telegram messages for Telegram-safe storage.
*/
+1 -1
View File
@@ -15,7 +15,7 @@ export interface S3FileRecord extends File {
* Repository interface for File entity persistence.
*
* Abstracts all file CRUD operations currently spread across
* `src/db/files.ts` and `src/db/files-ext.ts`.
* `infrastructure/persistence/repositories/file-repository.ts` and `infrastructure/persistence/repositories/file-repository.ts`.
*/
export interface IFileRepository {
/**
+1 -1
View File
@@ -3,7 +3,7 @@ import type { MultipartPart, MultipartUpload } from '../entities/multipart';
/**
* Repository interface for S3 multipart upload persistence.
*
* Abstracts the multipart upload operations currently in `src/db/multipart.ts`.
* Abstracts the multipart upload operations currently in `infrastructure/persistence/repositories/multipart-repository.ts`.
* Manages both multipart upload sessions and their individual parts.
*/
export interface IMultipartRepository {
+26 -2
View File
@@ -1,6 +1,7 @@
import { readFileSync } from 'node:fs';
import logger from './shared/logger/index';
import { TELEGRAM_CHUNK_SIZE_MAX_BYTES } from './shared/utils/validation';
import { PositiveIntSchema } from './shared/validation/schemas';
interface AppConfig {
/** Application version (read from package.json, kept in sync by semantic-release prepare.mjs) */
@@ -88,6 +89,27 @@ const parseNumber = (value: string | undefined, fallback: number): number => {
return Number.isFinite(parsed) && parsed > 0 ? parsed : fallback;
};
/**
* Parses the PORT env var with fail-fast validation.
*
* Defaults to 4000 ONLY when the variable is absent or empty. A present but
* invalid value ('abc', '-5', '0') throws so the service refuses to start
* misconfigured instead of silently listening on the wrong port.
*
* @param value - The raw `PORT` env var value.
* @returns The validated port number.
* @throws {Error} When the value is present but not a positive integer.
*/
const parsePort = (value: string | undefined): number => {
if (value === undefined || value === '') return 4000;
const parsed = PositiveIntSchema.safeParse(value);
if (!parsed.success) {
logger.error(`Invalid PORT value: ${JSON.stringify(value)} — PORT must be a positive integer`);
throw new Error('PORT must be a positive integer');
}
return parsed.data;
};
const parseTokens = (value: string | undefined): string[] =>
(value || '')
.split(',')
@@ -172,7 +194,7 @@ export const config: AppConfig = {
storageChatId: parseInt(process.env.STORAGE_CHANNEL_ID!, 10),
baseUrl: process.env.BASE_URL!,
databaseUrl: process.env.DATABASE_URL!,
port: parseInt(process.env.PORT!, 10) || 4000,
port: parsePort(process.env.PORT),
nodeEnv: process.env.NODE_ENV || 'development',
logLevel: process.env.LOG_LEVEL || 'info',
rateLimitWindowMs: parseNumber(process.env.RATE_LIMIT_WINDOW_MS, 60000),
@@ -197,7 +219,9 @@ export const config: AppConfig = {
),
};
logger.info('Environment variables loaded', {
// Debug-level: every import of env.ts would otherwise dump the full config
// (secrets masked, but still one noisy line per test file) to the log stream.
logger.debug('Environment variables loaded', {
config: {
...config,
botTokens: config.botTokens.map(maskSecret),
+32 -21
View File
@@ -49,28 +49,39 @@ const gracefulShutdown = async (signal: string): Promise<void> => {
process.on('SIGTERM', () => gracefulShutdown('SIGTERM'));
process.on('SIGINT', () => gracefulShutdown('SIGINT'));
// Periodic maintenance intervals
setInterval(cleanupRateLimitCache, 60000);
setInterval(
() => {
const removed = fileInfoCache.cleanup();
if (removed > 0) {
logger.info(`Cleaned up ${removed} expired cache entries`);
}
},
5 * 60 * 1000,
// Periodic maintenance intervals. Timers are unref'd so they never keep the
// process alive on their own (safe no-op where `unref` is unavailable).
const unref = (timer: unknown): void => {
if (typeof timer === 'object' && timer !== null && 'unref' in timer) {
(timer as { unref?: () => void }).unref?.();
}
};
unref(setInterval(cleanupRateLimitCache, 60000));
unref(
setInterval(
() => {
const removed = fileInfoCache.cleanup();
if (removed > 0) {
logger.info(`Cleaned up ${removed} expired cache entries`);
}
},
5 * 60 * 1000,
),
);
setInterval(
() => {
const snapshot = metricsCollector.getSnapshot();
logger.info('Metrics snapshot', {
uploadLatency: snapshot.uploadLatency,
uploadThroughput: snapshot.uploadThroughput.toFixed(2),
errorRate: snapshot.errorRate.toFixed(2),
cacheHitRate: snapshot.cacheHitRate.toFixed(2),
});
},
5 * 60 * 1000,
unref(
setInterval(
() => {
const snapshot = metricsCollector.getSnapshot();
logger.info('Metrics snapshot', {
uploadLatency: snapshot.uploadLatency,
uploadThroughput: snapshot.uploadThroughput.toFixed(2),
errorRate: snapshot.errorRate.toFixed(2),
cacheHitRate: snapshot.cacheHitRate.toFixed(2),
});
},
5 * 60 * 1000,
),
);
logger.info('Application running successfully');
@@ -21,7 +21,7 @@ const mapRowToBucket = (row: Record<string, unknown>): Bucket => ({
/**
* Drizzle-backed implementation of {@link IBucketRepository}.
*
* Delegates to the same SQL queries as the original `src/db/buckets.ts`
* Delegates to the same SQL queries as the original `infrastructure/persistence/repositories/bucket-repository.ts`
* module, using raw SQL for drizzle tables that are not part of the
* typed schema.
*/
@@ -31,7 +31,7 @@ const mapRowToFilePart = (row: Record<string, unknown>): FilePart => ({
/**
* Drizzle-backed implementation of {@link IFilePartRepository}.
*
* Delegates to the same SQL queries as the original `src/db/file-parts.ts`
* Delegates to the same SQL queries as the original `infrastructure/persistence/repositories/file-part-repository.ts`
* module, using raw SQL for all operations.
*/
export class DrizzleFilePartRepository implements IFilePartRepository {
@@ -51,8 +51,8 @@ const mapDbRowToS3Record = (row: Record<string, unknown>): S3FileRecord => ({
/**
* Drizzle-backed implementation of {@link IFileRepository}.
*
* Delegates to the same SQL queries as the original `src/db/files.ts` and
* `src/db/files-ext.ts` modules while presenting a clean domain interface.
* Delegates to the same SQL queries as the original `infrastructure/persistence/repositories/file-repository.ts` and
* `infrastructure/persistence/repositories/file-repository.ts` modules while presenting a clean domain interface.
*/
export class DrizzleFileRepository implements IFileRepository {
/**
@@ -19,7 +19,7 @@ const mapRowToMultipartUpload = (r: Record<string, unknown>): MultipartUpload =>
/**
* Drizzle-backed implementation of {@link IMultipartRepository}.
*
* Delegates to the same SQL queries as the original `src/db/multipart.ts`
* Delegates to the same SQL queries as the original `infrastructure/persistence/repositories/multipart-repository.ts`
* module, using raw SQL for all operations on the un-typed
* `multipart_uploads` and `multipart_parts` tables.
*/
@@ -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),
};
}),
);
+19
View File
@@ -0,0 +1,19 @@
/**
* Builds a Telegram CDN download URL from a file path and bot token.
*
* Single canonical implementation — previously constructed inline in
* `interfaces/http/controllers/file-controller.ts`,
* `interfaces/http/controllers/web-api-controller.ts`,
* `interfaces/http/controllers/s3-controller.ts`, and
* `infrastructure/telegram/chunked-storage.ts`.
*
* NOTE: URLs produced here embed the bot token. They must only be used
* server-side (outbound fetch to the Telegram CDN), never exposed to
* clients in redirects or response bodies.
*
* @param filePath - The Telegram file path returned by getFile.
* @param botToken - The bot token used to authenticate the download.
* @returns The full Telegram CDN URL.
*/
export const buildTelegramFileUrl = (filePath: string, botToken: string): string =>
`https://api.telegram.org/file/bot${botToken}/${filePath}`;
+6 -7
View File
@@ -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,12 +1,12 @@
import type { AuthSession } from '../../../application/dto/auth';
import {
type AuthSession,
createLoginUseCase,
createLogoutUseCase,
createMeUseCase,
} from '../../../application/use-cases/authenticate';
import { config } from '../../../env';
import { LoginBodySchema } from '../../../shared/validation/schemas';
import {
checkBearerToken,
clearSessionCookie,
createSessionCookie,
getAuthSession,
@@ -36,14 +36,17 @@ const notFound = (): Response => json({ error: 'Not found' }, 404);
/**
* Parses the login request body, extracting the `token` field.
*
* Validated through {@link LoginBodySchema} (the canonical boundary schema
* for `POST /api/v1/auth/login`).
*
* @param req - The incoming HTTP request with a JSON body.
* @returns The login token payload, or `null` when the body is invalid.
*/
const readLoginBody = async (req: Request): Promise<{ token: string } | null> => {
try {
const body = (await req.json()) as { token?: unknown };
if (typeof body.token !== 'string' || body.token.length === 0) return null;
return { token: body.token };
const parsed = LoginBodySchema.safeParse(await req.json());
if (!parsed.success) return null;
return { token: parsed.data.token };
} catch {
return null;
}
@@ -119,8 +122,10 @@ export const handleLogout = async (): Promise<Response> => {
export const handleMe = async (req: Request): Promise<Response> => {
if (!isAuthEnabled()) return notFound();
// getAuthSession already checks the bearer token (cookie first, then
// Authorization header) — no second check needed here.
const session: AuthSession | null = getAuthSession(req);
if (!session && !checkBearerToken(req.headers.get('authorization'))) {
if (!session) {
return json({ error: 'Unauthorized' }, 401);
}
@@ -132,13 +137,7 @@ export const handleMe = async (req: Request): Promise<Response> => {
},
});
const activeSession = session ?? {
username: 'admin',
expiresAt: null,
method: 'bearer' as const,
};
const result = await meUseCase(activeSession);
const result = await meUseCase(session);
if (!result) {
return json({ error: 'Unauthorized' }, 401);
@@ -1,9 +1,9 @@
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';
import { sanitizeFilenameHeader } from '../../../shared/http/filename';
import logger from '../../../shared/logger/index';
import { cleanupTempFile, formatCreatedAt, getErrorMessage } from '../../../shared/utils/file';
import { locateZipEntry } from '../../../shared/utils/zip';
@@ -19,53 +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;
};
/**
* Builds a Telegram CDN download URL from a file path and bot token.
*
* @param filePath - The Telegram file path returned by getFile.
* @param botToken - The bot token used to authenticate the download.
* @returns The full Telegram CDN URL.
*/
const buildTelegramFileUrl = (filePath: string, botToken: string): string =>
`https://api.telegram.org/file/bot${botToken}/${filePath}`;
/**
* Sanitises a file name for use in a Content-Disposition header, removing
* characters that could enable header injection.
*
* @param fileName - The raw file name.
* @returns The sanitised file name.
*/
const sanitizeFilenameHeader = (fileName: string): string =>
fileName.replace(/[\\"]/g, '').replace(/[\n\r]/g, '');
/**
* Returns a JSON error response with the given status code and message.
*
@@ -83,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.
@@ -117,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),
);
@@ -161,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.
File diff suppressed because one or more lines are too long
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,111 @@
import { bucketRepository, fileRepository } from '../../../../infrastructure/di';
import { BucketNameSchema, parseOrNull } from '../../../../shared/validation/schemas';
import { bucketVersioningConfigurationXml, listBucketsXml, s3ErrorResponse } from '../../../s3/xml';
import { resolveBucketOr404, s3Response } from './s3-common';
// ─────── Bucket Operations ───────
/**
* Handles GET / — lists all buckets as an S3 ListAllMyBuckets XML response.
*
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML response with the bucket list.
*/
export const handleListBuckets = async (reqId: string): Promise<Response> => {
const buckets = await bucketRepository.list();
const xml = listBucketsXml(buckets, reqId);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
/**
* Handles PUT /{bucket} — creates a new S3 bucket.
*
* Validates the bucket name format (BucketNameSchema, incl. the stricter M13
* rules: no consecutive dots, no IP format, no `xn--` prefix) and checks for
* duplicates.
*
* @param bucketName - The requested bucket name.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML response indicating success or failure.
*/
export const handleCreateBucket = async (bucketName: string, reqId: string): Promise<Response> => {
if (parseOrNull(BucketNameSchema, bucketName) === null) {
return s3ErrorResponse(
'InvalidBucketName',
'The specified bucket is not valid.',
`/${bucketName}`,
400,
reqId,
);
}
const existing = await bucketRepository.findByName(bucketName);
if (existing) {
return s3ErrorResponse(
'BucketAlreadyExists',
'The requested bucket name is not available.',
`/${bucketName}`,
409,
reqId,
);
}
await bucketRepository.create(bucketName);
return s3Response(null, 200, reqId);
};
/**
* Handles HEAD /{bucket} — checks whether a bucket exists.
*
* @param bucketName - The bucket name to check.
* @param reqId - The request identifier for S3 headers.
* @returns A 200 response when the bucket exists, or an S3 XML error.
*/
export const handleHeadBucket = async (bucketName: string, reqId: string): Promise<Response> => {
const bucket = await resolveBucketOr404(bucketRepository, bucketName, `/${bucketName}`, reqId);
if (bucket instanceof Response) return bucket;
return s3Response(null, 200, reqId);
};
/**
* Handles DELETE /{bucket} — deletes a bucket.
*
* Fails with `BucketNotEmpty` if the bucket still contains objects.
*
* @param bucketName - The bucket name to delete.
* @param reqId - The request identifier for S3 headers.
* @returns A 204 response on success, or an S3 XML error.
*/
export const handleDeleteBucket = async (bucketName: string, reqId: string): Promise<Response> => {
const bucket = await resolveBucketOr404(bucketRepository, bucketName, `/${bucketName}`, reqId);
if (bucket instanceof Response) return bucket;
const objCount = await fileRepository.countByBucket(bucket.id);
if (objCount > 0) {
return s3ErrorResponse(
'BucketNotEmpty',
'The bucket you tried to delete is not empty.',
`/${bucketName}`,
409,
reqId,
);
}
await bucketRepository.delete(bucketName);
return s3Response(null, 204, reqId);
};
/**
* Handles GET /{bucket}?versioning — returns the bucket versioning
* configuration (always disabled in this implementation).
*
* @param bucketName - The bucket name.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML response with the versioning configuration.
*/
export const handleGetBucketVersioning = async (
bucketName: string,
reqId: string,
): Promise<Response> => {
const bucket = await resolveBucketOr404(bucketRepository, bucketName, `/${bucketName}`, reqId);
if (bucket instanceof Response) return bucket;
return s3Response(bucketVersioningConfigurationXml(), 200, reqId, {
'content-type': 'application/xml',
});
};
@@ -0,0 +1,146 @@
import { nanoid } from 'nanoid';
import type { Bucket } from '../../../../domain/entities/bucket';
import type { MultipartUpload } from '../../../../domain/entities/multipart';
import type { IBucketRepository } from '../../../../domain/ports/bucket-repository';
import type { IMultipartRepository } from '../../../../domain/ports/multipart-repository';
import { config } from '../../../../env';
import { s3Headers } from '../../../s3/headers';
import { unsatisfiedContentRange } from '../../../s3/range';
import { s3ErrorResponse } from '../../../s3/xml';
/**
* The default S3 region returned when no region is explicitly configured.
*/
export const REGION = config.s3DefaultRegion || 'us-east-1';
/**
* Generates a unique request identifier for S3 responses.
*
* @returns A hex string suitable for x-amz-request-id and x-amz-id-2.
*/
export const REQUEST_ID = (): string => nanoid(16);
/**
* Builds a standard S3 response with the appropriate headers.
*
* @param body - The XML or empty response body.
* @param status - HTTP status code.
* @param reqId - The request identifier for S3 headers.
* @param extraHeaders - Optional extra response headers.
* @returns An S3-formatted Response.
*/
export const s3Response = (
body: string | null,
status: number,
reqId: string,
extraHeaders: Record<string, string> = {},
): Response => {
// Add content-type for empty 200-series responses (not 204 which has no body)
if (
body === null &&
status >= 200 &&
status < 300 &&
status !== 204 &&
!extraHeaders['content-type']
) {
extraHeaders['content-type'] = 'application/xml';
}
return new Response(body, { status, headers: s3Headers(reqId, extraHeaders) });
};
/**
* Resolves a bucket by name, returning a `NoSuchBucket` S3 error when missing.
*
* Replaces the ~10 identical `findByName` + `NoSuchBucket` blocks previously
* inlined in every bucket/object/listing/multipart handler.
*
* @param bucketRepo - The bucket repository to look up.
* @param bucket - The bucket name from the request path.
* @param path - The request path for the S3 error resource.
* @param reqId - The request identifier for S3 headers.
* @returns The bucket record, or an S3 error Response when not found.
*/
export const resolveBucketOr404 = async (
bucketRepo: IBucketRepository,
bucket: string,
path: string,
reqId: string,
): Promise<Bucket | Response> => {
const bucketRecord = await bucketRepo.findByName(bucket);
if (!bucketRecord) {
return s3ErrorResponse(
'NoSuchBucket',
'The specified bucket does not exist.',
path,
404,
reqId,
);
}
return bucketRecord;
};
/**
* Resolves an in-progress multipart upload, returning a `NoSuchUpload` S3
* error when missing.
*
* Replaces the 5 identical `findById` + `NoSuchUpload` blocks previously
* inlined in the multipart handlers. Pass `key` for the H5 key-match check
* (UploadPart / CompleteMultipartUpload); omit it for Abort / ListParts,
* which historically only checked existence.
*
* @param multipartRepo - The multipart repository to look up.
* @param uploadId - The upload identifier from `?uploadId=`.
* @param path - The request path for the S3 error resource.
* @param reqId - The request identifier for S3 headers.
* @param key - Optional object key the upload must belong to.
* @returns The upload record, or an S3 error Response when not found.
*/
export const requireUploadOr404 = async (
multipartRepo: IMultipartRepository,
uploadId: string,
path: string,
reqId: string,
key?: string,
): Promise<MultipartUpload | Response> => {
const multipart = await multipartRepo.findById(uploadId);
if (!multipart || (key !== undefined && multipart.s3Key !== key)) {
return s3ErrorResponse(
'NoSuchUpload',
'The specified upload does not exist.',
path,
404,
reqId,
);
}
return multipart;
};
/**
* Returns the stable S3 etag for a file, falling back to a random ID when
* the record has no content hash yet.
*
* Replaces the 4 inline `file.fileHash || nanoid(16)` fallbacks (conditional
* headers, HeadObject, copy result, list entries). Callers add quotes where
* the transport needs them. The copy-source check keeps its own inline
* `telegramFileId` variant (M9) with a comment at the call site.
*
* @param fileHash - The stored SHA-256 content hash (may be null).
* @returns The hash, or a random 16-char fallback.
*/
export const etagOrFallback = (fileHash: string | null): string => fileHash || nanoid(16);
/**
* Builds the shared 416 response for unsatisfiable Range requests.
*
* Replaces the 3 identical `InvalidRange` blocks (chunked GET, regular GET,
* multipart GET).
*
* @param path - The request path for the S3 error resource.
* @param totalSize - The total object size for the Content-Range header.
* @param reqId - The request identifier for S3 headers.
* @returns A 416 S3 error Response.
*/
export const invalidRangeResponse = (path: string, totalSize: number, reqId: string): Response =>
s3ErrorResponse('InvalidRange', 'The requested range is not satisfiable.', path, 416, reqId, {
'content-range': unsatisfiedContentRange(totalSize),
});
@@ -0,0 +1,142 @@
import type { File as FileEntity } from '../../../../domain/entities/file';
import { bucketRepository, fileRepository } from '../../../../infrastructure/di';
import { clampMaxKeys } from '../../../../shared/validation/schemas';
import { listBucketResultXml, listBucketV2ResultXml } from '../../../s3/xml';
import { etagOrFallback, resolveBucketOr404, s3Response } from './s3-common';
/** Shape of an S3 list entry object. */
export type S3ListEntry = {
key: string;
sizeBytes: number;
etag: string;
lastModified: Date;
mimeType: string;
};
/**
* Maps a File entity to an S3 list entry object.
*
* @param file - The file entity from the repository.
* @returns An S3 list entry with key, size, etag, last modified, and MIME type.
*/
export const mapFileToListEntry = (file: FileEntity): S3ListEntry => ({
key: file.s3Key ?? '',
sizeBytes: file.sizeBytes,
etag: etagOrFallback(file.fileHash),
lastModified: file.createdAt instanceof Date ? file.createdAt : new Date(),
mimeType: file.mimeType,
});
/**
* Handles GET /{bucket} (ListObjectsV1 with query parameters).
*
* @param bucket - The bucket name.
* @param searchParams - URL query parameters (prefix, delimiter, max-keys,
* marker, encoding-type).
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML ListBucketResult response.
*/
export const handleListObjectsV1 = async (
bucket: string,
searchParams: URLSearchParams,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(bucketRepository, bucket, `/${bucket}`, reqId);
if (bucketRecord instanceof Response) return bucketRecord;
const prefix = searchParams.get('prefix') || '';
const delimiter = searchParams.get('delimiter') || null;
const maxKeys = clampMaxKeys(searchParams.get('max-keys'));
const marker = searchParams.get('marker') || null;
const encodingType = searchParams.get('encoding-type') || null;
const { objects, prefixes: commonPrefixes } = await fileRepository.listByPrefix(
bucketRecord.id,
prefix,
delimiter,
maxKeys,
marker,
);
const isTruncated = objects.length > maxKeys;
const displayObjects = objects.slice(0, maxKeys);
const nextMarker = isTruncated
? (displayObjects[displayObjects.length - 1]?.s3Key ?? null)
: null;
const xml = listBucketResultXml(
bucket,
displayObjects.map(mapFileToListEntry),
commonPrefixes,
isTruncated,
marker,
maxKeys,
prefix,
delimiter,
nextMarker,
reqId,
encodingType,
);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
/**
* Handles GET /{bucket}?list-type=2 (ListObjectsV2).
*
* Small approved behavior fix: V2 previously used a bare
* `Math.min(parse, 1000)` which admitted 0, negatives, and NaN. It now uses
* the same `clampMaxKeys` ([1, 1000]) as V1.
*
* @param bucket - The bucket name.
* @param searchParams - URL query parameters (prefix, delimiter, max-keys,
* continuation-token, start-after, encoding-type).
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML ListBucketV2Result response.
*/
export const handleListObjectsV2 = async (
bucket: string,
searchParams: URLSearchParams,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(bucketRepository, bucket, `/${bucket}`, reqId);
if (bucketRecord instanceof Response) return bucketRecord;
const prefix = searchParams.get('prefix') || '';
const delimiter = searchParams.get('delimiter') || null;
const maxKeys = clampMaxKeys(searchParams.get('max-keys'));
const continuationToken = searchParams.get('continuation-token') || null;
const startAfter = searchParams.get('start-after') || null;
const encodingType = searchParams.get('encoding-type') || null;
const { objects, prefixes: commonPrefixes } = await fileRepository.listByPrefix(
bucketRecord.id,
prefix,
delimiter,
maxKeys,
continuationToken || startAfter,
);
const isTruncated = objects.length > maxKeys;
const displayObjects = objects.slice(0, maxKeys);
const nextContinuationToken = isTruncated
? (displayObjects[displayObjects.length - 1]?.s3Key ?? null)
: null;
const xml = listBucketV2ResultXml(
bucket,
displayObjects.map(mapFileToListEntry),
commonPrefixes,
isTruncated,
maxKeys,
prefix,
delimiter,
continuationToken,
nextContinuationToken,
displayObjects.length,
reqId,
encodingType,
);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
@@ -0,0 +1,431 @@
import { createReadStream } from 'node:fs';
import { nanoid } from 'nanoid';
import { buildNewFile } from '../../../../domain/entities/file-factory';
import type { IMultipartRepository } from '../../../../domain/ports/multipart-repository';
import type { ForwardResult } from '../../../../domain/ports/telegram-service';
import { config } from '../../../../env';
import {
bucketRepository,
fileRepository,
multipartRepository,
} from '../../../../infrastructure/di';
import { botPool } from '../../../../infrastructure/telegram/bot-pool';
import { cleanupTempFile, DEFAULT_FILE_TYPE } from '../../../../shared/utils/file';
import {
CompletePartSchema,
clampMaxKeys,
PartNumberSchema,
parseOrNull,
} from '../../../../shared/validation/schemas';
import {
completeMultipartUploadXml,
initiateMultipartUploadXml,
listMultipartUploadsXml,
listPartsXml,
parseCompleteMultipartBody,
s3ErrorResponse,
} from '../../../s3/xml';
import { requireUploadOr404, resolveBucketOr404, s3Response } from './s3-common';
type MultipartRepo = IMultipartRepository;
// ─────── Multipart Upload ───────
/**
* Handles POST /{bucket}/{key}?uploads — initiates a multipart upload.
*
* @param bucket - The bucket name.
* @param key - The object key being uploaded.
* @param _searchParams - URL query parameters (unused).
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML InitiateMultipartUpload response.
*/
export const handleCreateMultipartUpload = async (
bucket: string,
key: string,
_searchParams: URLSearchParams,
headers: Record<string, string>,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(
bucketRepository,
bucket,
`/${bucket}/${key}`,
reqId,
);
if (bucketRecord instanceof Response) return bucketRecord;
// NOTE: IMultipartRepository.create takes 3 args (bucketId, s3Key,
// initiatedBy). The content-type is intentionally dropped here — the
// complete step also falls back to 'application/octet-stream' since the
// in-progress upload record carries no contentType field.
void headers;
const uploadId = await (multipartRepository as MultipartRepo).create(bucketRecord.id, key, 's3');
const xml = initiateMultipartUploadXml(bucket, key, uploadId);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
/**
* Handles PUT /{bucket}/{key}?uploadId=&partNumber= — uploads a single
* part of a multipart upload.
*
* @param bucket - The bucket name.
* @param key - The object key.
* @param searchParams - URL query parameters containing uploadId and
* partNumber.
* @param req - The incoming HTTP request with the part body.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 response with the part etag, or an error.
*/
export const handleUploadPart = async (
bucket: string,
key: string,
searchParams: URLSearchParams,
req: Request,
reqId: string,
): Promise<Response> => {
const uploadId = searchParams.get('uploadId')!;
// M14: partNumber must be an integer 1–10000 (PartNumberSchema mirrors the
// old manual check; same InvalidArgument message preserved).
const partNumber = parseOrNull(PartNumberSchema, searchParams.get('partNumber'));
if (partNumber === null) {
return s3ErrorResponse(
'InvalidArgument',
'Part number must be an integer between 1 and 10000',
`/${bucket}/${key}`,
400,
reqId,
);
}
// H5: Verify both upload exists AND key matches.
const multipart = await requireUploadOr404(
multipartRepository,
uploadId,
`/${bucket}/${key}`,
reqId,
key,
);
if (multipart instanceof Response) return multipart;
// Stream the part body to temp — O(1) memory, safe for large parts
const tempPath = `/tmp/filedrop-mp-${nanoid()}`;
const writer = Bun.file(tempPath).writer();
const reader = (
req.body ??
new ReadableStream({
start(c) {
c.close();
},
})
).getReader();
const hasher = new Bun.CryptoHasher('sha256');
let sizeBytes = 0;
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
const chunk = Buffer.from(value);
sizeBytes += chunk.byteLength;
hasher.update(chunk);
writer.write(chunk);
}
await writer.end();
} catch (error) {
try {
writer.end();
} catch {
// ignore during error path
}
await cleanupTempFile(tempPath);
throw error;
} finally {
reader.releaseLock();
}
if (sizeBytes > config.telegramChunkSizeBytes) {
await cleanupTempFile(tempPath);
return s3ErrorResponse(
'EntityTooLarge',
`Your proposed upload part size (${sizeBytes} bytes) exceeds the maximum allowed part size (${config.telegramChunkSizeBytes} bytes) for this storage backend. Use smaller part sizes.`,
`/${bucket}/${key}`,
400,
reqId,
);
}
// M4: Ensure temp file cleanup even if forwardToStorage fails
let forwardResult: ForwardResult;
try {
forwardResult = await botPool.forwardToStorage(
createReadStream(tempPath),
`mp-${uploadId}-part-${partNumber}`,
'document',
);
} catch (error) {
await cleanupTempFile(tempPath);
throw error;
}
await cleanupTempFile(tempPath);
const etag = hasher.digest('hex');
await multipartRepository.insertPart({
uploadId,
partNumber,
telegramFileId: forwardResult.telegramFileId,
telegramFileUniqueId: forwardResult.telegramFileUniqueId,
storageMessageId: forwardResult.storageMessageId,
sizeBytes,
etag,
});
return s3Response(null, 200, reqId, { etag: `"${etag}"` });
};
/**
* Handles POST /{bucket}/{key}?uploadId= — completes a multipart upload.
*
* Validates the submitted part list (each part against CompletePartSchema —
* invalid → InvalidPart 400; all parts present, ascending order), creates
* the final file record, and marks the upload as completed.
*
* @param bucket - The bucket name.
* @param key - The object key.
* @param searchParams - URL query parameters containing uploadId.
* @param body - The raw XML request body containing the complete part list.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML CompleteMultipartUpload response.
*/
export const handleCompleteMultipartUpload = async (
bucket: string,
key: string,
searchParams: URLSearchParams,
body: string,
reqId: string,
): Promise<Response> => {
const uploadId = searchParams.get('uploadId')!;
// H5: Verify both upload exists AND key matches (consistent with handleUploadPart)
const multipart = await requireUploadOr404(
multipartRepository,
uploadId,
`/${bucket}/${key}`,
reqId,
key,
);
if (multipart instanceof Response) return multipart;
const rawParts = parseCompleteMultipartBody(body);
// Validate each submitted part against CompletePartSchema; any invalid
// part → InvalidPart 400 (approved small fix; previously malformed parts
// were silently dropped by the regex parser and surfaced as count mismatch).
const parts: Array<{ partNumber: number; etag: string }> = [];
for (const raw of rawParts) {
const valid = parseOrNull(CompletePartSchema, raw);
if (valid === null) {
return s3ErrorResponse(
'InvalidPart',
'One or more specified parts could not be found.',
`/${bucket}/${key}`,
400,
reqId,
);
}
parts.push(valid);
}
const storedParts = await multipartRepository.listParts(uploadId);
// Validate ascending part order
const partNumbers = parts.map((p) => p.partNumber);
if (partNumbers.length > 1 && partNumbers.some((n, i) => i > 0 && n <= partNumbers[i - 1])) {
return s3ErrorResponse(
'InvalidPartOrder',
'The list of parts was not in ascending order.',
`/${bucket}/${key}`,
400,
reqId,
);
}
// H8: Verify count AND part numbers AND etags match stored parts
if (parts.length !== storedParts.length) {
return s3ErrorResponse(
'InvalidPart',
'One or more specified parts could not be found.',
`/${bucket}/${key}`,
400,
reqId,
);
}
// Build a map for O(1) part number lookup
const storedByNumber = new Map<number, (typeof storedParts)[0]>();
for (const sp of storedParts) {
storedByNumber.set(sp.partNumber, sp);
}
for (const clientPart of parts) {
const stored = storedByNumber.get(clientPart.partNumber);
if (!stored || stored.etag !== clientPart.etag) {
return s3ErrorResponse(
'InvalidPart',
'One or more specified parts could not be found. The etag or part number does not match.',
`/${bucket}/${key}`,
400,
reqId,
);
}
}
const totalSize = storedParts.reduce((sum, p) => sum + Number(p.sizeBytes), 0);
const combinedEtag = storedParts.map((p) => p.etag).join('-');
const publicId = nanoid();
// The upload record carries no content-type, so the assembled object
// defaults to application/octet-stream.
const mimeType = 'application/octet-stream';
await fileRepository.create(
buildNewFile({
publicId,
telegramFileId: storedParts[0]!.telegramFileId,
telegramFileUniqueId: storedParts[0]!.telegramFileUniqueId,
storageChatId: config.storageChatId,
storageMessageId: storedParts[0]!.storageMessageId,
fileName: key.split('/').pop() || 'file',
mimeType,
sizeBytes: totalSize,
fileType: DEFAULT_FILE_TYPE,
uploaderId: 0,
bucketId: multipart.bucketId,
s3Key: key,
storageBackend: 'telegram',
multipartUploadId: uploadId,
}),
);
await multipartRepository.complete(uploadId);
const location = `${config.baseUrl}/${bucket}/${key}`;
const xml = completeMultipartUploadXml(bucket, key, combinedEtag, location);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
/**
* Handles GET /{bucket}?uploads — lists in-progress multipart uploads.
*
* @param bucket - The bucket name.
* @param searchParams - URL query parameters (max-uploads, key-marker).
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML ListMultipartUploadsResult response.
*/
export const handleListMultipartUploads = async (
bucket: string,
searchParams: URLSearchParams,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(bucketRepository, bucket, `/${bucket}`, reqId);
if (bucketRecord instanceof Response) return bucketRecord;
const maxUploads = clampMaxKeys(searchParams.get('max-uploads'));
const keyMarker = searchParams.get('key-marker') || null;
const { uploads, isTruncated, nextKeyMarker } = await multipartRepository.listByBucket(
bucketRecord.id,
maxUploads,
keyMarker,
);
const xml = listMultipartUploadsXml(
bucket,
uploads.map((u) => ({
key: u.s3Key,
uploadId: u.uploadId,
initiatedAt: u.initiatedAt,
initiatedBy: u.initiatedBy,
})),
maxUploads,
isTruncated,
nextKeyMarker,
reqId,
);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
/**
* Handles DELETE /{bucket}/{key}?uploadId= — aborts a multipart upload.
*
* @param bucket - The bucket name.
* @param key - The object key.
* @param searchParams - URL query parameters containing uploadId.
* @param reqId - The request identifier for S3 headers.
* @returns A 204 response on success, or an S3 XML error.
*/
export const handleAbortMultipartUpload = async (
bucket: string,
key: string,
searchParams: URLSearchParams,
reqId: string,
): Promise<Response> => {
const uploadId = searchParams.get('uploadId')!;
const multipart = await requireUploadOr404(
multipartRepository,
uploadId,
`/${bucket}/${key}`,
reqId,
);
if (multipart instanceof Response) return multipart;
await multipartRepository.abort(uploadId);
return s3Response(null, 204, reqId);
};
/**
* Handles GET /{bucket}/{key}?uploadId= — lists uploaded parts of a
* multipart upload.
*
* @param bucket - The bucket name.
* @param key - The object key.
* @param searchParams - URL query parameters containing uploadId and
* optional max-parts.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML ListPartsResult response.
*/
export const handleListParts = async (
bucket: string,
key: string,
searchParams: URLSearchParams,
reqId: string,
): Promise<Response> => {
const uploadId = searchParams.get('uploadId')!;
const multipart = await requireUploadOr404(
multipartRepository,
uploadId,
`/${bucket}/${key}`,
reqId,
);
if (multipart instanceof Response) return multipart;
const parts = await multipartRepository.listParts(uploadId);
const maxParts = clampMaxKeys(searchParams.get('max-parts'));
const xml = listPartsXml(
bucket,
key,
uploadId,
parts.map((p) => ({
partNumber: p.partNumber,
etag: p.etag,
sizeBytes: p.sizeBytes,
createdAt: p.createdAt,
})),
maxParts,
false,
reqId,
);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
@@ -0,0 +1,358 @@
import type { File as FileEntity } from '../../../../domain/entities/file';
import {
bucketRepository,
chunkedStorage,
fileRepository,
multipartRepository,
} from '../../../../infrastructure/di';
import { botPool } from '../../../../infrastructure/telegram/bot-pool';
import { buildTelegramFileUrl } from '../../../../infrastructure/telegram/file-url';
import logger from '../../../../shared/logger/index';
import { getErrorMessage } from '../../../../shared/utils/file';
import { s3Headers } from '../../../s3/headers';
import { createGetObjectResponse, type ObjectPartSource } from '../../../s3/object-stream';
import { parseRangeHeader } from '../../../s3/range';
import { s3ErrorResponse } from '../../../s3/xml';
import { etagOrFallback, invalidRangeResponse, resolveBucketOr404, s3Response } from './s3-common';
// ─────── Conditional Headers Helper ──────
/**
* S3-compatible response for 304 Not Modified.
*/
export const notModifiedResponse = (
reqId: string,
etag: string,
mimeType: string,
sizeBytes: number,
lastModified: Date,
): Response =>
new Response(null, {
status: 304,
headers: s3Headers(reqId, {
etag,
'content-type': mimeType,
'content-length': String(sizeBytes),
'last-modified': lastModified.toUTCString(),
'x-amz-version-id': 'null',
}),
});
/**
* S3-compatible response for 412 Precondition Failed.
*/
export const preconditionFailedResponse = (path: string, reqId: string): Response =>
s3ErrorResponse(
'PreconditionFailed',
'At least one of the pre-conditions you specified did not hold.',
path,
412,
reqId,
);
/**
* Checks conditional headers (If-Match, If-None-Match, If-Modified-Since,
* If-Unmodified-Since) and returns a prepared Response if the condition
* is not satisfied, or `null` to let the request proceed.
*
* @returns A 304 / 412 Response when a condition fails, or `null` to continue.
*/
export const checkConditionalHeaders = (
headers: Record<string, string>,
file: {
mimeType: string;
sizeBytes: number;
fileHash: string | null;
createdAt: Date | string | number;
},
path: string,
reqId: string,
): Response | null => {
const etag = `"${etagOrFallback(file.fileHash)}"`;
const lastModified = file.createdAt instanceof Date ? file.createdAt : new Date(file.createdAt);
// If-Match
const ifMatch = headers['if-match'];
if (ifMatch && ifMatch !== '*' && ifMatch !== etag) {
return preconditionFailedResponse(path, reqId);
}
// If-None-Match
const ifNoneMatch = headers['if-none-match'];
if (ifNoneMatch && ifNoneMatch === etag) {
return notModifiedResponse(reqId, etag, file.mimeType, file.sizeBytes, lastModified);
}
// If-Modified-Since
const ifModifiedSince = headers['if-modified-since'];
if (ifModifiedSince) {
const since = new Date(ifModifiedSince);
if (!Number.isNaN(since.getTime()) && lastModified.getTime() <= since.getTime()) {
return notModifiedResponse(reqId, etag, file.mimeType, file.sizeBytes, lastModified);
}
}
// If-Unmodified-Since
const ifUnmodifiedSince = headers['if-unmodified-since'];
if (ifUnmodifiedSince) {
const since = new Date(ifUnmodifiedSince);
if (!Number.isNaN(since.getTime()) && lastModified.getTime() > since.getTime()) {
return preconditionFailedResponse(path, reqId);
}
}
return null;
};
// ─────── Object Operations ───────
/**
* Handles GET /{bucket}/{key} — retrieves an S3 object.
*
* Supports chunked objects (streaming multi-part response), multipart
* objects (assembled from a completed multipart upload), and regular
* Telegram-stored objects (proxy streaming or 302 redirect depending
* on configuration). HTTP Range headers are respected when present.
*
* @param bucket - The bucket name.
* @param key - The object key.
* @param _searchParams - URL query parameters (unused for GET).
* @param headers - The request headers (used for Range and etag checks).
* @param reqId - The request identifier for S3 headers.
* @returns An S3 response with the object content or an error.
*/
export const handleGetObject = async (
bucket: string,
key: string,
_searchParams: URLSearchParams,
headers: Record<string, string>,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(
bucketRepository,
bucket,
`/${bucket}/${key}`,
reqId,
);
if (bucketRecord instanceof Response) return bucketRecord;
const file = await fileRepository.findByBucketAndKey(bucketRecord.id, key);
if (!file)
return s3ErrorResponse(
'NoSuchKey',
'The specified key does not exist.',
`/${bucket}/${key}`,
404,
reqId,
);
// H3: Conditional headers — If-Match / If-None-Match / If-Modified-Since / If-Unmodified-Since
const conditionResult = checkConditionalHeaders(headers, file, `/${bucket}/${key}`, reqId);
if (conditionResult) {
return conditionResult;
}
// Chunked storage object
if (file.storageBackend === 'chunked') {
const totalSize = Number(file.sizeBytes);
const range = parseRangeHeader(headers.range || null, totalSize);
if (range.type === 'invalid') {
return invalidRangeResponse(`/${bucket}/${key}`, totalSize, reqId);
}
try {
return await chunkedStorage.createChunkedObjectResponse({ file, range, reqId });
} catch (error) {
logger.warn('Chunked object content fetch failed', { key, error: getErrorMessage(error) });
return s3ErrorResponse(
'InternalError',
'Failed to fetch object content from storage',
`/${bucket}/${key}`,
502,
reqId,
);
}
}
// Multipart upload assembled object
if (file.multipartUploadId) {
return handleGetMultipartObject(file, bucket, key, headers, reqId);
}
// Regular Telegram object
const fileInfo = await botPool.getFileInfo(file.telegramFileId);
const telegramUrl = buildTelegramFileUrl(fileInfo.file_path, fileInfo.bot_token);
const totalSize = file.sizeBytes;
const range = parseRangeHeader(headers.range || null, totalSize);
if (range.type === 'invalid') {
return invalidRangeResponse(`/${bucket}/${key}`, totalSize, reqId);
}
// H1: Always proxy S3 GETs to avoid leaking the Telegram bot token
// in redirect URLs. The 302 redirect path is removed because the
// URL contains the bot_token — exposing it to clients is a security risk.
const part: ObjectPartSource = {
telegramFileId: file.telegramFileId,
telegramUrl,
sizeBytes: file.sizeBytes,
partNumber: 1,
};
try {
return await createGetObjectResponse({
reqId,
contentType: file.mimeType,
etag: file.fileHash || '',
lastModified: file.createdAt instanceof Date ? file.createdAt : new Date(file.createdAt),
totalSize: file.sizeBytes,
parts: [part],
range,
});
} catch (error) {
logger.warn('Telegram content fetch failed', {
fileId: file.telegramFileId,
error: getErrorMessage(error),
});
return s3ErrorResponse(
'InternalError',
'Failed to fetch object content from storage',
`/${bucket}/${key}`,
502,
reqId,
);
}
};
/**
* Handles GET for objects assembled from a completed multipart upload.
*
* Resolves the Telegram CDN URLs for each part and builds a multi-part
* streaming response, respecting HTTP Range headers.
*
* @param file - The file entity with a `multipartUploadId` reference.
* @param bucket - The bucket name.
* @param key - The object key.
* @param headers - The request headers (for Range parsing).
* @param reqId - The request identifier for S3 headers.
* @returns An S3 response streaming the assembled object content.
*/
export const handleGetMultipartObject = async (
file: FileEntity,
bucket: string,
key: string,
headers: Record<string, string>,
reqId: string,
): Promise<Response> => {
const uploadId = file.multipartUploadId!;
const parts = await multipartRepository.listParts(uploadId);
if (parts.length === 0) {
return s3ErrorResponse(
'InternalError',
'Multipart object has no parts.',
`/${bucket}/${key}`,
500,
reqId,
);
}
const totalSize = parts.reduce((sum, p) => sum + Number(p.sizeBytes), 0);
const range = parseRangeHeader(headers.range || null, totalSize);
if (range.type === 'invalid') {
return invalidRangeResponse(`/${bucket}/${key}`, totalSize, reqId);
}
const sources: ObjectPartSource[] = [];
// Resolve all part CDN URLs concurrently (independent getFile calls) so
// assembly latency is ~1 round-trip instead of N.
sources.push(
...(await Promise.all(
parts.map(async (part) => {
const fileInfo = await botPool.getFileInfo(part.telegramFileId);
return {
telegramFileId: part.telegramFileId,
telegramUrl: buildTelegramFileUrl(fileInfo.file_path, fileInfo.bot_token),
sizeBytes: part.sizeBytes,
partNumber: part.partNumber,
};
}),
)),
);
// H1: Always proxy — never expose bot token in redirect URL
try {
return await createGetObjectResponse({
reqId,
contentType: file.mimeType,
etag: file.fileHash || parts.map((p) => p.etag).join('-'),
lastModified: file.createdAt instanceof Date ? file.createdAt : new Date(file.createdAt),
totalSize,
parts: sources,
range,
});
} catch (error) {
logger.warn('Telegram multipart content fetch failed', {
uploadId: file.multipartUploadId,
error: getErrorMessage(error),
});
return s3ErrorResponse(
'InternalError',
'Failed to fetch object content from storage',
`/${bucket}/${key}`,
502,
reqId,
);
}
};
/**
* Handles HEAD /{bucket}/{key} — returns object metadata without the body.
*
* @param bucket - The bucket name.
* @param key - The object key.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 response with object metadata headers.
*/
export const handleHeadObject = async (
bucket: string,
key: string,
headers: Record<string, string>,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(
bucketRepository,
bucket,
`/${bucket}/${key}`,
reqId,
);
if (bucketRecord instanceof Response) return bucketRecord;
const file = await fileRepository.findByBucketAndKey(bucketRecord.id, key);
if (!file)
return s3ErrorResponse(
'NoSuchKey',
'The specified key does not exist.',
`/${bucket}/${key}`,
404,
reqId,
);
// H3: Conditional headers for HEAD — If-Match / If-None-Match / If-Modified-Since / If-Unmodified-Since
const headConditionResult = checkConditionalHeaders(headers, file, `/${bucket}/${key}`, reqId);
if (headConditionResult) {
return headConditionResult;
}
return s3Response(null, 200, reqId, {
'content-type': file.mimeType,
'content-length': String(file.sizeBytes),
etag: `"${etagOrFallback(file.fileHash)}"`,
'last-modified':
file.createdAt instanceof Date ? file.createdAt.toUTCString() : new Date().toUTCString(),
'accept-ranges': 'bytes',
'cache-control': 'public, max-age=31536000',
'x-amz-version-id': 'null',
});
};
@@ -0,0 +1,442 @@
import { createReadStream } from 'node:fs';
import { nanoid } from 'nanoid';
import { buildNewFile } from '../../../../domain/entities/file-factory';
import type { ForwardResult } from '../../../../domain/ports/telegram-service';
import { config } from '../../../../env';
import { bucketRepository, chunkedStorage, fileRepository } from '../../../../infrastructure/di';
import { botPool } from '../../../../infrastructure/telegram/bot-pool';
import { cleanupTempFile, DEFAULT_FILE_TYPE, ensureExtension } from '../../../../shared/utils/file';
import { streamToTemp } from '../../../../shared/utils/temp-stream';
import { DeleteObjectsBodySchema, parseOrNull } from '../../../../shared/validation/schemas';
import { verifyBodyHash } from '../../../s3/auth';
import {
copyObjectResultXml,
deleteResultXml,
parseDeleteObjectsBody,
s3ErrorResponse,
} from '../../../s3/xml';
import { etagOrFallback, resolveBucketOr404, s3Response } from './s3-common';
/**
* Streams the request body to a temporary file while computing its SHA-256
* and MD5 hashes.
*
* Unlike `req.arrayBuffer()`, this approach uses O(1) memory regardless of
* file size, making it safe for multi-GB Docker registry layer blobs.
*
* MD5 is computed alongside SHA-256 so that Content-MD5 verification (when
* the header is present) does not need to re-read the entire file.
*
* @param body - The ReadableStream from the HTTP request body.
* @returns The temp file path, SHA-256 hash, MD5 hash (base64), total size, and signature bytes.
*/
export const streamBodyToTemp = async (
body: ReadableStream<Uint8Array> | null,
): Promise<{
tempPath: string;
fileHash: string;
md5Hash?: string;
sizeBytes: number;
signatureBuffer: Buffer;
}> => {
const reader = (
body ??
new ReadableStream({
start(c) {
c.close();
},
})
).getReader() as ReadableStreamDefaultReader<Uint8Array>;
return streamToTemp(reader, { computeMd5: true, prefix: '/tmp/filedrop-s3-' });
};
/**
* Handles PUT /{bucket}/{key} — uploads an S3 object.
*
* Streams the request body directly to a temporary file to avoid buffering
* the entire payload in memory. This is essential for supporting large
* Docker registry layer blobs (100MB–2GB+).
*
* Supports regular binary uploads, copy-object via `x-amz-copy-source`,
* and tag operations. Large files are stored as chunked objects (across
* multiple Telegram messages), while smaller files use a single Telegram
* message.
*
* @param bucket - The bucket name.
* @param key - The object key.
* @param searchParams - URL query parameters.
* @param headers - The request headers.
* @param req - The incoming HTTP request with the object body.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 response with the object etag or an error.
*/
export const handlePutObject = async (
bucket: string,
key: string,
searchParams: URLSearchParams,
headers: Record<string, string>,
req: Request,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(
bucketRepository,
bucket,
`/${bucket}/${key}`,
reqId,
);
if (bucketRecord instanceof Response) return bucketRecord;
// Tag operations are idempotent no-ops
if (searchParams.has('tagging')) {
return s3Response(null, 204, reqId);
}
// Copy-object path
const copySource = headers['x-amz-copy-source'];
if (copySource) {
return handleCopyObject(bucket, key, copySource, headers, bucketRecord.id, reqId);
}
// Stream body to temp file — O(1) memory, safe for multi-GB blobs
const contentType = headers['content-type'] || 'application/octet-stream';
const streamed = await streamBodyToTemp(req.body);
// H4: Verify body hash against x-amz-content-sha256
const bodyHashError = verifyBodyHash(streamed.fileHash, headers);
if (bodyHashError) {
await cleanupTempFile(streamed.tempPath);
return s3ErrorResponse(
bodyHashError.errorCode || 'BadDigest',
'The x-amz-content-sha256 you specified did not match what we received.',
`/${bucket}/${key}`,
400,
reqId,
);
}
// Content-Length validation: ensure actual body size matches header
const contentLengthHeader = headers['content-length'];
if (contentLengthHeader) {
const declaredLength = Number.parseInt(contentLengthHeader, 10);
if (Number.isFinite(declaredLength) && declaredLength !== streamed.sizeBytes) {
await cleanupTempFile(streamed.tempPath);
return s3ErrorResponse(
'IncompleteBody',
'You did not provide the number of bytes specified by the Content-Length HTTP header.',
`/${bucket}/${key}`,
400,
reqId,
);
}
}
// Content-MD5 validation: use pre-computed MD5 from streaming (no OOM re-read)
const contentMd5 = headers['content-md5'];
if (contentMd5 && contentMd5 !== streamed.md5Hash) {
await cleanupTempFile(streamed.tempPath);
return s3ErrorResponse(
'BadDigest',
'The Content-MD5 you specified did not match what we received.',
`/${bucket}/${key}`,
400,
reqId,
);
}
// M12: Reject oversized bodies
if (streamed.sizeBytes > config.maxRequestBodyBytes) {
await cleanupTempFile(streamed.tempPath);
return s3ErrorResponse(
'EntityTooLarge',
'Your proposed upload exceeds the maximum allowed object size.',
`/${bucket}/${key}`,
400,
reqId,
);
}
// Idempotent PUT: if the object already exists, skip upload
try {
const existing = await fileRepository.findByBucketAndKey(bucketRecord.id, key);
if (existing) {
await cleanupTempFile(streamed.tempPath);
return s3Response(null, 200, reqId, { etag: `"${streamed.fileHash}"` });
}
return await storeFileFromTemp(streamed, key, bucketRecord, contentType, reqId);
} catch (error) {
await cleanupTempFile(streamed.tempPath);
throw error;
}
};
/**
* Stores a streamed file to Telegram storage as an S3 object.
*
* Accepts the result of `streamBodyToTemp` (temp path + hash + size) instead
* of a raw Buffer, enabling O(1) memory usage for multi-GB Docker layer blobs.
*
* Handles both chunked (large files) and single-message (small files) paths.
*
* @param streamed - The streamed file result (temp path, hash, size, signature).
* @param key - The S3 object key.
* @param bucketRecord - The resolved bucket record (id and name).
* @param contentType - The MIME type from the request Content-Type header.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 response with the etag of the stored object.
*/
export const storeFileFromTemp = async (
streamed: { tempPath: string; fileHash: string; sizeBytes: number; signatureBuffer: Buffer },
key: string,
bucketRecord: { id: string; name: string },
contentType: string,
reqId: string,
): Promise<Response> => {
const fileName = key.split('/').pop() || 'file';
const { fileName: finalFileName, mimeType } = ensureExtension(
fileName,
streamed.signatureBuffer,
contentType,
);
const bucketId = bucketRecord.id;
const partFileNamePrefix = `s3-${bucketRecord.name}-${key.replace(/\//g, '_')}`;
if (streamed.sizeBytes > config.telegramChunkSizeBytes) {
const file = await chunkedStorage.storeFileInTelegramChunks({
tempPath: streamed.tempPath,
partFileNamePrefix,
fileName: finalFileName,
mimeType,
sizeBytes: streamed.sizeBytes,
fileType: DEFAULT_FILE_TYPE,
uploaderId: 0,
bucketId,
s3Key: key,
});
await cleanupTempFile(streamed.tempPath);
return s3Response(null, 200, reqId, { etag: `"${file.fileHash}"` });
}
const fileStream = createReadStream(streamed.tempPath);
let forwardResult: ForwardResult;
try {
forwardResult = await botPool.forwardToStorage(fileStream, partFileNamePrefix, 'document');
} catch (error) {
fileStream.destroy();
throw error;
}
fileStream.destroy();
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,
s3Key: key,
storageBackend: 'telegram',
}),
);
await cleanupTempFile(streamed.tempPath);
return s3Response(null, 200, reqId, { etag: `"${streamed.fileHash}"` });
};
/**
* Handles PUT /{bucket}/{key} with an `x-amz-copy-source` header.
*
* Creates a new file record referencing the same Telegram-stored data as
* the source object. Chunked source objects are not supported for copy.
*
* @param _destBucket - The destination bucket name (unused — bucket record
* already resolved).
* @param destKey - The destination object key.
* @param rawCopySource - The raw `x-amz-copy-source` header value.
* @param headers - The request headers (for conditional copy checks).
* @param destBucketId - The UUID of the destination bucket.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML response with the copy result or an error.
*/
export const handleCopyObject = async (
_destBucket: string,
destKey: string,
rawCopySource: string,
headers: Record<string, string>,
destBucketId: string,
reqId: string,
): Promise<Response> => {
const copySource = decodeURIComponent(rawCopySource);
const sourcePath = copySource.startsWith('/') ? copySource.slice(1) : copySource;
const parts = sourcePath.split('/');
const sourceBucket = parts[0];
const sourceKey = parts.slice(1).join('/');
const sourceBucketRecord = await bucketRepository.findByName(sourceBucket);
if (!sourceBucketRecord)
return s3ErrorResponse(
'NoSuchBucket',
'The specified bucket does not exist.',
copySource,
404,
reqId,
);
const sourceFile = await fileRepository.findByBucketAndKey(sourceBucketRecord.id, sourceKey);
if (!sourceFile)
return s3ErrorResponse(
'NoSuchKey',
'The specified key does not exist.',
copySource,
404,
reqId,
);
// Chunked objects cannot be copied yet
if (sourceFile.storageBackend === 'chunked') {
return s3ErrorResponse(
'NotImplemented',
'Copying chunked objects is not yet implemented.',
copySource,
501,
reqId,
);
}
// Conditional copy: if-match / if-none-match checks
// M9: Use stable etag (telegramFileId fallback when fileHash is null) —
// kept inline (not etagOrFallback) because this variant falls back to
// telegramFileId while the read/list paths fall back to nanoid(16).
const sourceEtag = sourceFile.fileHash || sourceFile.telegramFileId;
const ifMatch = headers['x-amz-copy-source-if-match'];
const ifNoneMatch = headers['x-amz-copy-source-if-none-match'];
if (ifMatch && ifMatch !== '*' && ifMatch !== `"${sourceEtag}"`) {
return s3ErrorResponse(
'PreconditionFailed',
'The preconditions you specified did not hold.',
copySource,
412,
reqId,
);
}
if (ifNoneMatch && ifNoneMatch === `"${sourceEtag}"`) {
return s3ErrorResponse(
'PreconditionFailed',
'The preconditions you specified did not hold.',
copySource,
412,
reqId,
);
}
const publicId = nanoid();
await fileRepository.create(
buildNewFile({
publicId,
telegramFileId: sourceFile.telegramFileId,
telegramFileUniqueId: sourceFile.telegramFileUniqueId,
storageChatId: sourceFile.storageChatId,
storageMessageId: sourceFile.storageMessageId,
fileName: sourceFile.fileName,
mimeType: sourceFile.mimeType,
sizeBytes: sourceFile.sizeBytes,
fileType: sourceFile.fileType,
uploaderId: 0,
fileHash: sourceFile.fileHash,
bucketId: destBucketId,
s3Key: destKey,
storageBackend: 'telegram',
}),
);
// copyObjectResultXml quotes the etag itself.
const xml = copyObjectResultXml(etagOrFallback(sourceFile.fileHash), new Date());
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
/**
* Handles DELETE /{bucket}/{key} — soft-deletes an S3 object.
*
* @param bucket - The bucket name.
* @param key - The object key to delete.
* @param reqId - The request identifier for S3 headers.
* @returns A 204 response on success, or an S3 XML error.
*/
export const handleDeleteObject = async (
bucket: string,
key: string,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(
bucketRepository,
bucket,
`/${bucket}/${key}`,
reqId,
);
if (bucketRecord instanceof Response) return bucketRecord;
await fileRepository.softDelete(bucketRecord.id, key);
return s3Response(null, 204, reqId);
};
/**
* Handles POST /{bucket}?delete — batch-deletes multiple S3 objects.
*
* Parses the XML Delete request body, soft-deletes each key, and returns
* an XML delete result. The parsed body is validated with
* DeleteObjectsBodySchema (malformed output → MalformedXML 400; the schema
* also enforces the M11 S3 limit of 1000 keys).
*
* @param bucket - The bucket name.
* @param body - The raw XML request body.
* @param reqId - The request identifier for S3 headers.
* @returns An S3 XML response listing deleted keys.
*/
export const handleDeleteObjects = async (
bucket: string,
body: string,
reqId: string,
): Promise<Response> => {
const bucketRecord = await resolveBucketOr404(bucketRepository, bucket, `/${bucket}`, reqId);
if (bucketRecord instanceof Response) return bucketRecord;
const parsed = parseOrNull(DeleteObjectsBodySchema, parseDeleteObjectsBody(body));
// M11: S3 spec limits batch delete to 1000 keys (also enforced by schema)
if (parsed === null || parsed.keys.length > 1000) {
return s3ErrorResponse(
'MalformedXML',
'The XML you provided was not well-formed or did not validate against our published schema. Max 1000 keys per request.',
`/${bucket}`,
400,
reqId,
);
}
const { keys, quiet } = parsed;
const deletedKeys: string[] = [];
const errors: Array<{ key: string; code: string; message: string }> = [];
for (const key of keys) {
const ok = await fileRepository.softDelete(bucketRecord.id, key);
if (ok) {
deletedKeys.push(key);
} else {
// Per S3 spec, deleting a non-existent key is idempotent — report as success
deletedKeys.push(key);
}
}
const xml = quiet ? deleteResultXml([], []) : deleteResultXml(deletedKeys, errors);
return s3Response(xml, 200, reqId, { 'content-type': 'application/xml' });
};
@@ -0,0 +1,238 @@
import { config } from '../../../../env';
import logger from '../../../../shared/logger/index';
import { getErrorMessage } from '../../../../shared/utils/file';
import { verifyPresignedUrl, verifySignature } from '../../../s3/auth';
import { S3_CORS_HEADERS } from '../../../s3/headers';
import { s3ErrorResponse } from '../../../s3/xml';
import {
handleCreateBucket,
handleDeleteBucket,
handleGetBucketVersioning,
handleHeadBucket,
handleListBuckets,
} from './s3-bucket-handlers';
import { REGION, REQUEST_ID, s3Response } from './s3-common';
import { handleListObjectsV1, handleListObjectsV2 } from './s3-listing';
import {
handleAbortMultipartUpload,
handleCompleteMultipartUpload,
handleCreateMultipartUpload,
handleListMultipartUploads,
handleListParts,
handleUploadPart,
} from './s3-multipart-handlers';
import { handleGetObject, handleHeadObject } from './s3-object-read';
import { handleDeleteObject, handleDeleteObjects, handlePutObject } from './s3-object-write';
/**
* Builds an S3 OPTIONS preflight response with CORS headers.
*
* @returns A 204 No Content Response.
*/
export const s3OptionsResponse = (): Response =>
new Response(null, { status: 204, headers: S3_CORS_HEADERS });
/**
* Parses an S3 pathname into bucket and key components.
*
* Supports path-style URLs such as `/bucket-name/key/with/prefix`.
*
* @param pathname - The URL pathname.
* @returns An object with the extracted bucket and key (both may be null).
*/
export const parseS3Path = (pathname: string): { bucket: string | null; key: string | null } => {
const parts = pathname.split('/').filter(Boolean);
if (parts.length === 0) return { bucket: null, key: null };
if (parts.length === 1) return { bucket: parts[0], key: null };
// Decode URI components to match virtual-hosted behavior (H10)
const key = parts
.slice(1)
.map((segment) => decodeURIComponent(segment))
.join('/');
return { bucket: parts[0], key };
};
/**
* Converts a Request's headers into a plain key-value record (all keys
* lowercased) for SigV4 signature verification.
*
* @param req - The incoming HTTP request.
* @returns A record of lowercased header key-value pairs.
*/
export const headersToRecord = (req: Request): Record<string, string> => {
const record: Record<string, string> = {};
for (const [key, value] of req.headers.entries()) {
record[key.toLowerCase()] = value;
}
return record;
};
/**
* Main S3 request dispatcher.
*
* Parses the request (method, path, query parameters, headers), validates
* the SigV4 signature or presigned URL, and dispatches to the appropriate
* bucket, object, or multipart operation handler.
*
* Supports both path-style (`/bucket/key`) and virtual-hosted-style
* (`bucket.example.com/key`) addressing.
*
* @param req - The incoming S3 HTTP request.
* @param virtualHostBucket - When the request was routed through a
* virtual-hosted domain, the extracted bucket
* name; otherwise `null`.
* @returns An S3-formatted Response.
*/
export const handleS3Request = async (
req: Request,
virtualHostBucket: string | null = null,
): Promise<Response> => {
const method = req.method;
const url = new URL(req.url);
const pathname = url.pathname;
const { bucket, key } = virtualHostBucket
? {
bucket: virtualHostBucket,
key: pathname === '/' ? null : decodeURIComponent(pathname.slice(1)),
}
: parseS3Path(pathname);
const headers = headersToRecord(req);
const searchParams = url.searchParams;
const reqId = REQUEST_ID();
// Handle CORS preflight
if (method === 'OPTIONS') {
return s3OptionsResponse();
}
// SigV4 authentication
const isPresigned = searchParams.has('X-Amz-Signature');
const authResult = isPresigned
? await verifyPresignedUrl({
url: req.url,
method,
headers,
s3AccessKey: config.s3AccessKey,
s3SecretKey: config.s3SecretKey,
region: REGION,
})
: await verifySignature(
method,
req.url,
headers,
null,
config.s3AccessKey,
config.s3SecretKey,
REGION,
);
if (!authResult.isValid) {
const status = authResult.errorCode === 'NotImplemented' ? 501 : 403;
const message =
authResult.errorCode === 'NotImplemented'
? 'aws-chunked streaming payloads are not supported.'
: isPresigned
? 'Presigned URL verification failed'
: 'Authentication required';
return s3ErrorResponse(
authResult.errorCode || 'AccessDenied',
message,
pathname,
status,
reqId,
);
}
try {
// ── Root: ListBuckets / Service-level operations ──
if (!bucket) {
if (method === 'GET') {
return handleListBuckets(reqId);
}
return s3ErrorResponse(
'MethodNotAllowed',
'The specified method is not allowed against this resource.',
'/',
405,
reqId,
);
}
// ── Bucket-level operations ──
if (!key) {
if (method === 'GET') {
if (searchParams.has('versioning')) {
return handleGetBucketVersioning(bucket, reqId);
}
if (searchParams.has('uploads')) {
return handleListMultipartUploads(bucket, searchParams, reqId);
}
const listType = searchParams.get('list-type');
if (listType === '2') {
return handleListObjectsV2(bucket, searchParams, reqId);
}
return handleListObjectsV1(bucket, searchParams, reqId);
}
if (method === 'PUT') return handleCreateBucket(bucket, reqId);
if (method === 'HEAD') return handleHeadBucket(bucket, reqId);
if (method === 'DELETE') return handleDeleteBucket(bucket, reqId);
if (method === 'POST') {
if (searchParams.has('delete')) {
const body = await req.text();
return handleDeleteObjects(bucket, body, reqId);
}
if (searchParams.has('tagging')) {
return s3Response(null, 204, reqId);
}
}
return s3ErrorResponse(
'MethodNotAllowed',
'The specified method is not allowed against this resource.',
`/${bucket}`,
405,
reqId,
);
}
// ── Object-level: Multipart operations ──
if (searchParams.has('uploads') && method === 'POST') {
return handleCreateMultipartUpload(bucket, key, searchParams, headers, reqId);
}
if (searchParams.has('uploadId') && searchParams.has('partNumber') && method === 'PUT') {
return handleUploadPart(bucket, key, searchParams, req, reqId);
}
if (searchParams.has('uploadId') && method === 'POST') {
const body = await req.text();
return handleCompleteMultipartUpload(bucket, key, searchParams, body, reqId);
}
if (searchParams.has('uploadId') && method === 'DELETE') {
return handleAbortMultipartUpload(bucket, key, searchParams, reqId);
}
if (searchParams.has('uploadId') && method === 'GET') {
return handleListParts(bucket, key, searchParams, reqId);
}
// ── Standard object operations ──
if (method === 'GET') return handleGetObject(bucket, key, searchParams, headers, reqId);
if (method === 'HEAD') return handleHeadObject(bucket, key, headers, reqId);
if (method === 'PUT') return handlePutObject(bucket, key, searchParams, headers, req, reqId);
if (method === 'DELETE') return handleDeleteObject(bucket, key, reqId);
return s3ErrorResponse(
'MethodNotAllowed',
'The specified method is not allowed against this resource.',
`/${bucket}/${key}`,
405,
reqId,
);
} catch (error: unknown) {
logger.error('S3 operation error', { bucket, key, error: getErrorMessage(error) });
return s3ErrorResponse(
'InternalError',
'We encountered an internal error. Please try again.',
pathname,
500,
reqId,
);
}
};
@@ -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,17 +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';
import { BucketNameSchema } from '../../../shared/validation/schemas';
/** 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.
@@ -66,7 +81,7 @@ export const handleListBucketsV1 = async (): Promise<Response> => {
*/
export const handleCreateBucketV1 = async (req: Request): Promise<Response> => {
const body = (await req.json()) as { name?: string };
if (!body.name || !/^[a-z0-9][a-z0-9.-]{1,61}[a-z0-9]$/.test(body.name)) {
if (!body.name || !BucketNameSchema.safeParse(body.name).success) {
return jsonError('Invalid bucket name. Use lowercase, 3-63 chars, no underscore', 400);
}
const existing = await bucketRepository.findByName(body.name);
@@ -170,73 +185,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 +229,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 +254,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 },
);
}
};
/**
+4 -30
View File
@@ -1,5 +1,6 @@
import { createHmac, timingSafeEqual } from 'node:crypto';
import { createHmac } from 'node:crypto';
import { config } from '../../../env';
import { timingSafeCompare } from '../../../shared/utils/crypto';
const ADMIN_USERNAME = 'admin';
const SIGNATURE_SEPARATOR = '.';
@@ -11,17 +12,7 @@ type Handler = (req: Request) => Response | Promise<Response>;
* Represents an authenticated user session after successful
* authentication via cookie or bearer token.
*/
export interface AuthSession {
/** The authenticated username (always "admin" in this implementation). */
username: string;
/**
* Expiration date of the session, or `null` for bearer-token
* sessions which do not expire at the session level.
*/
expiresAt: Date | null;
/** The authentication method used to establish this session. */
method: 'cookie' | 'bearer';
}
export type { AuthSession } from '../../../application/dto/auth';
/** Options for configuring cookie-based session behaviour. */
interface CookieOptions {
@@ -64,24 +55,7 @@ const decodePayload = (value: string): string | null => {
*/
export const isAuthEnabled = (secret = config.adminApiToken): boolean => secret.length > 0;
/**
* Compares two strings using a timing-safe algorithm to prevent
* timing side-channel attacks.
*
* @param left - First string to compare.
* @param right - Second string to compare.
* @returns `true` when the strings are equal, `false` otherwise.
*/
export 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);
};
export { timingSafeCompare } from '../../../shared/utils/crypto';
/**
* Signs an arbitrary payload string with HMAC-SHA256 using the given
+51 -8
View File
@@ -1,4 +1,3 @@
import { handleSwaggerHtml, handleSwaggerJson } from '../../../routes/swagger';
import { getS3RouteBucket, shouldHandleS3 } from '../../../shared/utils/s3-detection';
import { handleLogin, handleLogout, handleMe } from '../controllers/auth-controller';
import { handleFileInfo, handleFileRedirect } from '../controllers/file-controller';
@@ -9,6 +8,7 @@ import { handleUpload } from '../controllers/upload-controller';
import { handleWebApiV1 } from '../controllers/web-api-controller';
import { requireAuth } from '../middleware/auth';
import { withRateLimit } from '../middleware/rate-limit';
import { handleSwaggerHtml, handleSwaggerJson } from '../swagger';
/**
* Dispatches an S3 request directly, bypassing rate limiting.
@@ -16,7 +16,7 @@ import { withRateLimit } from '../middleware/rate-limit';
* S3 API calls (used by Docker registry for blob pushes) must not be
* rate-limited — large concurrent layer uploads would hit the limit and
* fail. The Docker registry client retries on 5xx, not 4xx, so a 429
* would abort the entire push.
* would abort the entire push. Do NOT wrap this in withRateLimit.
*
* @param req - The incoming S3 request.
* @returns The S3 response.
@@ -25,6 +25,40 @@ const handleS3Direct = (req: Request): Promise<Response> => {
return handleS3Request(req, getS3RouteBucket(req));
};
/**
* Generic CORS preflight for non-S3 API requests.
*
* S3 preflights are answered by the S3 controller (S3 XML CORS headers);
* anything else (dashboard fetches, future API endpoints) gets a plain
* permissive 204 so browsers can proceed.
*
* @returns A 204 No Content Response with permissive CORS headers.
*/
const apiOptionsResponse = (): Response =>
new Response(null, {
status: 204,
headers: {
'Access-Control-Allow-Origin': '*',
'Access-Control-Allow-Methods': 'GET, PUT, HEAD, DELETE, POST, PATCH, OPTIONS',
'Access-Control-Allow-Headers':
'Authorization, Content-Type, X-Amz-Date, X-Amz-Content-Sha256',
},
});
/**
* Handles an OPTIONS request on the catch-all route.
*
* S3 clients preflight with SigV4 headers — those go to the S3 handler.
* Anything else is a generic API preflight and gets a plain 204.
*
* @param req - The incoming OPTIONS request.
* @returns The S3 or generic CORS preflight response.
*/
const handleCatchAllOptions = (req: Request): Promise<Response> => {
if (shouldHandleS3(req, Object.fromEntries(req.headers))) return handleS3Direct(req);
return Promise.resolve(apiOptionsResponse());
};
/**
* Defines all HTTP routes for the application.
*
@@ -58,7 +92,7 @@ export const routes = {
},
'/': {
GET: (req: Request): Promise<Response> => {
if (shouldHandleS3(req)) return handleS3Direct(req);
if (shouldHandleS3(req, Object.fromEntries(req.headers))) return handleS3Direct(req);
return handleHome();
},
PUT: (req: Request): Promise<Response> => {
@@ -66,10 +100,19 @@ export const routes = {
if (shouldHandleS3(req, headers)) return handleS3Direct(req);
return Promise.resolve(new Response('Not Allowed', { status: 405 }));
},
HEAD: handleS3Direct,
DELETE: handleS3Direct,
POST: handleS3Direct,
OPTIONS: handleS3Direct,
HEAD: (req: Request): Promise<Response> => {
if (shouldHandleS3(req, Object.fromEntries(req.headers))) return handleS3Direct(req);
return Promise.resolve(new Response('Not Found', { status: 404 }));
},
DELETE: (req: Request): Promise<Response> => {
if (shouldHandleS3(req, Object.fromEntries(req.headers))) return handleS3Direct(req);
return Promise.resolve(new Response('Not Found', { status: 404 }));
},
POST: (req: Request): Promise<Response> => {
if (shouldHandleS3(req, Object.fromEntries(req.headers))) return handleS3Direct(req);
return Promise.resolve(new Response('Not Found', { status: 404 }));
},
OPTIONS: handleCatchAllOptions,
},
// Catch-all for S3 path-style requests (/{bucket}/{key} ...)
// Only intercepts requests with S3 auth headers; others get 404.
@@ -98,7 +141,7 @@ export const routes = {
if (shouldHandleS3(req, Object.fromEntries(req.headers))) return handleS3Direct(req);
return Promise.resolve(new Response('Not Found', { status: 404 }));
},
OPTIONS: handleS3Direct,
OPTIONS: handleCatchAllOptions,
},
'/api/v1/auth/login': {
POST: withRateLimit(handleLogin),
@@ -1,4 +1,4 @@
import { config } from '../env';
import { config } from '../../env';
const errorSchema = (example: string) => ({
type: 'object',
@@ -50,8 +50,9 @@ export const handleSwaggerJson = async (): Promise<Response> => {
openapi: '3.0.0',
info: {
title: 'FileDrop API',
version: '1.0.0',
description: 'File upload API with stream-based downloads.',
version: config.appVersion,
description:
'File upload API with stream-based downloads, S3-compatible object storage, and admin auth.',
},
servers: [
{
@@ -204,6 +205,156 @@ export const handleSwaggerJson = async (): Promise<Response> => {
},
},
},
'/api/v1/auth/login': {
post: {
summary: 'Admin Login',
description:
'Validates the admin API token and sets a signed session cookie. Returns 404 when auth is disabled.',
requestBody: {
required: true,
content: jsonContent(
objectSchema({
token: { type: 'string', example: 'admin-secret-token' },
}),
),
},
responses: {
'200': {
description: 'Login successful; session cookie set.',
content: jsonContent(
objectSchema({
username: { type: 'string', example: 'admin' },
}),
),
},
'400': {
description: 'Token is required.',
content: jsonContent(errorSchema('Token is required')),
},
'401': {
description: 'Invalid token.',
content: jsonContent(errorSchema('Invalid token')),
},
},
},
},
'/api/v1/auth/logout': {
post: {
summary: 'Admin Logout',
description: 'Clears the session cookie.',
responses: {
'200': {
description: 'Logout successful.',
content: jsonContent(
objectSchema({
success: { type: 'boolean', example: true },
}),
),
},
},
},
},
'/api/v1/auth/me': {
get: {
summary: 'Current User',
description:
'Returns the authenticated user from the session cookie or bearer token. Returns 404 when auth is disabled.',
responses: {
'200': {
description: 'User info.',
content: jsonContent(
objectSchema({
username: { type: 'string', example: 'admin' },
expiresAt: {
type: 'string',
format: 'date-time',
nullable: true,
example: '2026-05-18T10:00:00.000Z',
},
}),
),
},
'401': {
description: 'Unauthorized.',
content: jsonContent(errorSchema('Unauthorized')),
},
},
},
},
'/api/v1/{path}': {
get: {
summary: 'Web API (read)',
description:
'Public read endpoints: list buckets/objects and download files. See the dashboard for the full reference.',
responses: {
'200': {
description: 'Requested resource.',
},
'404': {
description: 'Not found.',
content: jsonContent(errorSchema('Not found')),
},
},
},
},
'/{bucket}': {
get: {
summary: 'S3 Bucket Operations',
description:
'S3-compatible bucket endpoint (SigV4 auth). Supports ListObjects, versioning queries, and bucket management. Served without rate limiting so Docker registry pushes are not aborted by 429s.',
parameters: [
{
name: 'bucket',
in: 'path',
required: true,
description: 'Bucket name.',
schema: { type: 'string' },
},
],
responses: {
'200': {
description: 'S3 XML response.',
},
'403': {
description: 'Signature mismatch.',
},
},
},
},
'/{bucket}/{key}': {
get: {
summary: 'S3 Object Operations',
description:
'S3-compatible object endpoint (SigV4 auth): GetObject, PutObject, DeleteObject, and multipart uploads. Served without rate limiting so Docker registry pushes are not aborted by 429s.',
parameters: [
{
name: 'bucket',
in: 'path',
required: true,
description: 'Bucket name.',
schema: { type: 'string' },
},
{
name: 'key',
in: 'path',
required: true,
description: 'Object key.',
schema: { type: 'string' },
},
],
responses: {
'200': {
description: 'S3 XML or object bytes.',
},
'403': {
description: 'Signature mismatch.',
},
'404': {
description: 'NoSuchBucket / NoSuchKey.',
},
},
},
},
},
};
+1 -23
View File
@@ -1,26 +1,4 @@
import { timingSafeEqual } from 'node:crypto';
/**
* Timing-safe string comparison that prevents timing attacks.
*
* Uses `crypto.timingSafeEqual` which runs in constant time regardless of
* where the strings differ. Returns false for mismatched-length inputs
* to avoid leaking length information via early return.
*
* @param left - The first string to compare.
* @param right - The second string to compare.
* @returns True if both strings are equal.
*/
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);
};
import { timingSafeCompare } from '../../shared/utils/crypto';
export interface SigV4Result {
isValid: boolean;
+7 -5
View File
@@ -7,11 +7,13 @@ const stripPort = (host: string): string => {
return host.split(':')[0].toLowerCase().replace(/\.$/, '');
};
const isValidBucketLabel = (bucket: string): boolean =>
/^[a-z0-9][a-z0-9.-]{1,61}[a-z0-9]$/.test(bucket) &&
!bucket.includes('..') &&
!bucket.includes('.-') &&
!bucket.includes('-.');
import { BucketNameSchema } from '../../shared/validation/schemas';
/**
* Validates a virtual-hosted bucket label against the single canonical
* bucket-name schema (same rules as bucket creation).
*/
const isValidBucketLabel = (bucket: string): boolean => BucketNameSchema.safeParse(bucket).success;
export const extractS3BucketFromHost = (host: string, domains: string[]): string | null => {
const normalizedHost = stripPort(host);
+13
View File
@@ -0,0 +1,13 @@
/**
* Sanitises a file name for use in a Content-Disposition header, removing
* characters that could enable header injection.
*
* Single canonical implementation — previously only present in
* `interfaces/http/controllers/file-controller.ts` while other download
* paths (S3 GET, web-api) did not sanitise at all.
*
* @param fileName - The raw file name.
* @returns The sanitised file name.
*/
export const sanitizeFilenameHeader = (fileName: string): string =>
fileName.replace(/[\\"]/g, '').replace(/[\n\r]/g, '');
+7 -2
View File
@@ -82,11 +82,16 @@ class MetricsCollector {
p95: this.calculatePercentile(this.uploadTimes, 95),
p99: this.calculatePercentile(this.uploadTimes, 99),
},
// Cumulative mean rate since process start (total requests / elapsed
// minutes). Intended for the coarse 5-minute ops log in index.ts, not
// a sliding-window throughput gauge.
uploadThroughput: this.totalRequests > 0 ? this.totalRequests / 60 : 0,
queueSize: 0, // Will be updated by queue
// No upload queue or bot-utilization tracker exists yet — both stay 0
// until one is wired in. Kept in the snapshot shape for compatibility.
queueSize: 0,
errorRate,
cacheHitRate,
botUtilization: 0, // Will be updated by bot tracker
botUtilization: 0,
timestamp: Date.now(),
};
}

Some files were not shown because too many files have changed in this diff Show More