Commit 9aa9709a by Xianquan Committed by GitHub

fix(storage): avoid MinIO XML entity expansion limit (#7356)

* fix(storage): avoid MinIO XML entity expansion limit

* fix(storage): add MinIO delete request timeout

* fix(storage): handle MinIO delete failures and timeouts

* test(storage): add shared adapter integration suite

* test(storage): expand adapter boundary coverage

* test(storage): recreate stable integration buckets

* fix(storage): align object key and delete contracts

* docs(storage): remove integration test design

* fix(storage): stabilize integration test setup

* fix(storage): keep legacy test mock import path

---------

Co-authored-by: Archer <545436317@qq.com>
parent 43bd4f9c
......@@ -26,3 +26,4 @@ FE_DOMAIN=https://fastgpt.example.com
1. Fixed an issue where Chatbox displayed system tool errors during streaming responses.
2. Fixed an issue where plain-text tool responses in full run details could be incorrectly parsed as Markdown, causing formatting issues.
3. Fixed an issue where switching the embedding model triggered training but did not rebuild vectors for existing data.
4. Fixed MinIO prefix-based bulk deletion failures caused by the XML entity expansion limit and added request timeout protection.
......@@ -24,3 +24,4 @@ FE_DOMAIN=https://fastgpt.example.com
1. chatbox 流输出时候,不应该展示系统工具的错误。
2. 完整运行详情,纯文本的工具响应 UI 可能会被 Markdown 错误解析,格式错乱。
3. 切换向量模型后,训练任务会触发但已有数据的向量未重建。
4. 修复 MinIO 按前缀批量删除大量对象时,可能因 XML 实体展开限制失败的问题,并增加请求超时保护。
import z from 'zod';
import { UploadConstraintsSchema } from '../contracts/type';
import { StorageObjectKeySchema, UploadConstraintsSchema } from '../contracts/type';
import { UploadFileHintSchema, UploadPolicySchema } from '../uploadPolicy/type';
import { S3_DOWNLOAD_URL_BATCH_MAX_SIZE } from '@fastgpt-sdk/storage/access-link';
......@@ -10,7 +10,7 @@ const HexSha256Schema = z
.regex(/^[a-f0-9]+$/);
export const S3AccessBucketNameSchema = z.string().min(1);
export const S3AccessObjectKeySchema = z.string().min(1);
export const S3AccessObjectKeySchema = StorageObjectKeySchema;
export const S3DownloadAliasIdSchema = UrlSafeTokenSchema.min(12).max(32);
export const S3DownloadAliasKeySchema = HexSha256Schema;
......
......@@ -5,6 +5,19 @@ import {
UploadFileHintSchema,
UploadPolicySchema
} from '../uploadPolicy/type';
import { assertStorageObjectKey } from '@fastgpt-sdk/storage';
/** FastGPT 入口与底层 Storage SDK 共用同一套对象 key 规范。 */
export const StorageObjectKeySchema = z.string().superRefine((key, context) => {
try {
assertStorageObjectKey(key);
} catch (error) {
context.addIssue({
code: z.ZodIssueCode.custom,
message: error instanceof Error ? error.message : 'Invalid storage object key'
});
}
});
export const S3MetadataSchema = z.object({
filename: z.string(),
......@@ -38,7 +51,7 @@ export type StorageDownloadUrlMode = z.infer<typeof StorageDownloadUrlModeSchema
export const CreatePostPresignedUrlParamsSchema = z.object({
filename: z.string().min(1),
rawKey: z.string().min(1),
rawKey: StorageObjectKeySchema,
contentType: UploadFileHintSchema.shape.contentType,
declaredExtension: UploadFileHintSchema.shape.declaredExtension,
declaredFilename: UploadFileHintSchema.shape.declaredFilename,
......@@ -64,7 +77,7 @@ export const CreatePostPresignedUrlResultSchema = z.object({
});
export type CreatePostPresignedUrlResult = z.infer<typeof CreatePostPresignedUrlResultSchema>;
export const CreateGetPresignedUrlParamsSchema = z.object({
key: z.string().nonempty(),
key: StorageObjectKeySchema,
expiredHours: z.number().positive().optional(),
responseContentType: z.string().nonempty().optional()
});
......@@ -74,7 +87,7 @@ export const UploadImage2S3BucketParamsSchema = z
.object({
base64Img: z.string().nonempty().optional(),
buffer: z.instanceof(Buffer).optional(),
uploadKey: z.string().nonempty(),
uploadKey: StorageObjectKeySchema,
mimetype: z.string().nonempty(),
filename: z.string().nonempty(),
expiredTime: z.coerce.date().optional()
......@@ -87,7 +100,7 @@ export type UploadImage2S3BucketParams = z.infer<typeof UploadImage2S3BucketPara
export const UploadFileByBodySchema = z.object({
body: z.union([z.instanceof(Buffer), z.string(), z.instanceof(Readable)]),
contentType: z.string().optional(),
key: z.string().nonempty(),
key: StorageObjectKeySchema,
filename: z.string().nonempty(),
expiredTime: z.coerce.date().optional()
});
......
......@@ -5,7 +5,11 @@ import type { Readable } from 'node:stream';
import { MongoS3TTL } from './models/ttl';
import { S3Buckets } from './config/constants';
import { S3PrivateBucket } from './buckets/private';
import { S3Sources, type UploadImage2S3BucketParams } from './contracts/type';
import {
S3Sources,
type UploadImage2S3BucketParams,
UploadImage2S3BucketParamsSchema
} from './contracts/type';
import { S3PublicBucket } from './buckets/public';
import { getNanoid } from '@fastgpt/global/common/string/tools';
import path from 'node:path';
......@@ -123,7 +127,14 @@ export async function uploadImage2S3Bucket(
bucketName: keyof typeof S3Buckets,
params: UploadImage2S3BucketParams
) {
const { base64Img, buffer: inputBuffer, filename, mimetype, uploadKey, expiredTime } = params;
const {
base64Img,
buffer: inputBuffer,
filename,
mimetype,
uploadKey,
expiredTime
} = UploadImage2S3BucketParamsSchema.parse(params);
const bucket = bucketName === 'private' ? new S3PrivateBucket() : new S3PublicBucket();
......
......@@ -59,4 +59,21 @@ describe('executeS3DeleteJob', () => {
})
).rejects.toThrow('Failed to delete 1 S3 object');
});
it('throws when multi-key deletion reports failed keys so BullMQ can retry', async () => {
global.s3BucketMap = {
'fastgpt-private': {
client: {
deleteObjectsByMultiKeys: vi.fn(async () => ({ keys: ['dataset/team/failed.txt'] }))
}
}
} as any;
await expect(
executeS3DeleteJob({
bucketName: 'fastgpt-private',
keys: ['dataset/team/deleted.txt', 'dataset/team/failed.txt']
})
).rejects.toThrow('Failed to delete 1 S3 object');
});
});
import { describe, expect, it } from 'vitest';
import {
CreateGetPresignedUrlParamsSchema,
CreatePostPresignedUrlParamsSchema,
StorageObjectKeySchema,
UploadFileByBodySchema,
UploadImage2S3BucketParamsSchema
} from '@fastgpt/service/common/s3/contracts/type';
import { S3AccessObjectKeySchema } from '@fastgpt/service/common/s3/accessLink/type';
type KeyParseResult =
| { success: true }
| { success: false; error: { issues: Array<{ message: string }> } };
const keySchemas: ReadonlyArray<readonly [string, (key: string) => KeyParseResult]> = [
['shared key', (key) => StorageObjectKeySchema.safeParse(key)],
[
'presigned PUT',
(key) => CreatePostPresignedUrlParamsSchema.safeParse({ filename: 'file.txt', rawKey: key })
],
['presigned GET', (key) => CreateGetPresignedUrlParamsSchema.safeParse({ key })],
[
'image upload',
(key) =>
UploadImage2S3BucketParamsSchema.safeParse({
uploadKey: key,
mimetype: 'image/png',
filename: 'file.png',
buffer: Buffer.from('image')
})
],
[
'body upload',
(key) => UploadFileByBodySchema.safeParse({ key, body: 'body', filename: 'file.txt' })
],
['access link', (key) => S3AccessObjectKeySchema.safeParse(key)]
];
describe('FastGPT storage object key schemas', () => {
it.each(keySchemas)('%s accepts portable URL-sensitive characters', (_name, parseKey) => {
expect(parseKey('team/folder # & + % ?/\u6587\u4ef6-\ud83d\ude00.txt').success).toBe(true);
});
it.each(keySchemas)('%s rejects a path the storage SDK would reject', (_name, parseKey) => {
const result = parseKey('team//file.txt');
expect(result.success).toBe(false);
if (!result.success) {
expect(result.error.issues[0]?.message).toContain('consecutive slashes');
}
});
it.each(keySchemas)('%s rejects a key beyond 850 UTF-8 bytes', (_name, parseKey) => {
expect(parseKey('a'.repeat(851)).success).toBe(false);
});
});
......@@ -5,10 +5,24 @@ import {
truncateFilename,
S3_FILENAME_MAX_LENGTH,
isS3ObjectKey,
getFileS3Key
getFileS3Key,
uploadImage2S3Bucket
} from '@fastgpt/service/common/s3/utils';
import * as stringTools from '@fastgpt/global/common/string/tools';
describe('uploadImage2S3Bucket', () => {
it('rejects an invalid key at the FastGPT boundary before uploading', async () => {
await expect(
uploadImage2S3Bucket('private', {
buffer: Buffer.from('image'),
filename: 'image.png',
mimetype: 'image/png',
uploadKey: 'dataset//image.png'
})
).rejects.toThrow('consecutive slashes');
});
});
describe('truncateFilename', () => {
it('should return filename as-is if within max length', () => {
const filename = 'short.pdf';
......
......@@ -38,6 +38,11 @@ const sdkMocks = vi.hoisted(() => ({
parsePkg: vi.fn()
}));
const loggerMocks = vi.hoisted(() => ({
warning: vi.fn(),
error: vi.fn()
}));
vi.mock('../../../src/service/mongo/models/tool', () => ({
MarketplaceToolIndexZodSchema: {
parse: (value: any) => ({
......@@ -66,10 +71,7 @@ vi.mock('../../../src/service/logger', () => ({
API: 'api'
}
},
getLogger: () => ({
warning: vi.fn(),
error: vi.fn()
})
getLogger: () => loggerMocks
}));
const createIndex = ({
......@@ -127,6 +129,8 @@ describe('PluginRepo', () => {
}
});
sdkMocks.parsePkg.mockReset();
loggerMocks.warning.mockReset();
loggerMocks.error.mockReset();
Reflect.deleteProperty(globalThis, 'marketplacePluginManifestCache');
});
......@@ -273,6 +277,29 @@ describe('PluginRepo', () => {
);
});
it('reports failed asset keys without failing a completed manifest publish', async () => {
const existing = createIndex({ pluginId: 'tool-a', version: '1.0.0', etag: 'old-etag' });
const record = createManifest(
createIndex({ pluginId: 'tool-a', version: '1.0.0', etag: 'new-etag' })
);
modelMocks.findOne.mockReturnValue(mockLean(existing));
modelMocks.updateOne.mockResolvedValue({ acknowledged: true });
s3Mocks.uploadJsonToS3.mockResolvedValue(undefined);
s3Mocks.deleteObjectsByPrefixFromS3.mockResolvedValue({
keys: ['assets/tool-a/1.0.0/old-etag/failed.svg']
});
const { PluginRepo } = await import('../../../src/service/plugin/repo');
await expect(new PluginRepo().publishToolManifest(record)).resolves.toBeUndefined();
expect(loggerMocks.warning).toHaveBeenCalledWith(
'Delete old marketplace tool assets partially failed',
expect.objectContaining({
failedKeys: ['assets/tool-a/1.0.0/old-etag/failed.svg']
})
);
});
it('rejects non-official manifests when the toolId already exists under another source', async () => {
const record = createManifest(
createIndex({
......
# Integration tests are opt-in. Copy the required provider block to .env.test.local.
# BUCKET must start with fastgpt-sdk-. Its existing contents are deleted before every suite.
# Local MinIO
STORAGE_TEST_MINIO_ENABLED=false
STORAGE_TEST_MINIO_BUCKET=fastgpt-sdk.integration-bucket-1
STORAGE_TEST_MINIO_ENDPOINT=http://127.0.0.1:9000
STORAGE_TEST_MINIO_REGION=us-east-1
STORAGE_TEST_MINIO_ACCESS_KEY_ID=minioadmin
STORAGE_TEST_MINIO_SECRET_ACCESS_KEY=minioadmin
# AWS S3 or another S3-compatible endpoint
STORAGE_TEST_AWS_S3_ENABLED=false
STORAGE_TEST_AWS_S3_BUCKET=
STORAGE_TEST_AWS_S3_ENDPOINT=https://s3.amazonaws.com
STORAGE_TEST_AWS_S3_REGION=us-east-1
STORAGE_TEST_AWS_S3_ACCESS_KEY_ID=
STORAGE_TEST_AWS_S3_SECRET_ACCESS_KEY=
STORAGE_TEST_AWS_S3_FORCE_PATH_STYLE=false
# Alibaba Cloud OSS
STORAGE_TEST_OSS_ENABLED=false
STORAGE_TEST_OSS_BUCKET=
STORAGE_TEST_OSS_ENDPOINT=
STORAGE_TEST_OSS_REGION=
STORAGE_TEST_OSS_ACCESS_KEY_ID=
STORAGE_TEST_OSS_SECRET_ACCESS_KEY=
# Tencent Cloud COS
STORAGE_TEST_COS_ENABLED=false
# Full COS bucket name; it must end with -<STORAGE_TEST_COS_APP_ID>.
STORAGE_TEST_COS_BUCKET=
STORAGE_TEST_COS_REGION=
STORAGE_TEST_COS_APP_ID=
STORAGE_TEST_COS_ACCESS_KEY_ID=
STORAGE_TEST_COS_SECRET_ACCESS_KEY=
.env
.env.local
.env.test
......@@ -148,14 +148,14 @@ const storage = createStorage({
> 重要:当前实现状态(以代码为准):
> - `generatePresignedPutUrl`:**AWS S3 / MinIO / COS / OSS 已实现**。
> - `generatePresignedGetUrl`:目前各 adapter 仍为 **未实现**(会抛 `Error('Method not implemented.')`)
> - `generatePresignedGetUrl`:**AWS S3 / MinIO / COS / OSS 已实现**
### 预签名 PUT 直传示例(浏览器 / 前端)
`generatePresignedPutUrl` 返回的 `metadata` 字段语义更接近“需要带上的 headers”(不同厂商前缀不同,如 `x-oss-meta-*` / `x-cos-meta-*`)。
```ts
const { putUrl, metadata } = await storage.generatePresignedPutUrl({
const { url: putUrl, metadata } = await storage.generatePresignedPutUrl({
key: 'demo/hello.txt',
expiredSeconds: 600,
metadata: { app: 'fastgpt', purpose: 'direct-upload' }
......@@ -179,11 +179,13 @@ await fetch(putUrl, {
- **`NoSuchBucketError`**: bucket 不存在(部分 adapter 会用它包装底层错误)。
- **`NoBucketReadPermissionError`**: bucket 无读取权限(部分 adapter 会用它包装底层错误)。
- **`EmptyObjectError`**: 下载时对象为空(例如底层 SDK 返回 `Body` 为空)。
- **`InvalidStorageObjectKeyError`**: key/prefix 未通过 SDK 统一预检;`reason``field``actualBytes``maxBytes` 可用于结构化处理。
建议你在业务层做分层处理:可恢复错误(重试/提示权限)与不可恢复错误(配置错误/接口未实现)。
## 注意事项
- **key 使用统一规范**:所有 adapter 都在远端请求前要求 1 - 850 UTF-8 bytes,拒绝前导 `/`、反斜线、连续 `//`、控制字符和 `.`/`..` 路径段;空格及 `+ # & % ?`、中文、emoji 可正常使用。
- **按前缀删除是高危操作**`prefix` 必须是非空字符串;强烈建议使用业务隔离前缀(例如 `team/{teamId}/`),避免误删整桶。
- **metadata 厂商差异**:不同厂商对元数据 key 前缀/大小写/可用字符/大小限制不同,建议使用简单 ASCII key,并控制总体大小。
- **流式下载/上传**:大文件建议使用 `Readable`,减少内存峰值。
......@@ -191,10 +193,29 @@ await fetch(putUrl, {
## 开发与构建
```bash
pnpm -C FastGPT/packages/storage dev
pnpm -C FastGPT/packages/storage build
pnpm --filter @fastgpt-sdk/storage dev
pnpm --filter @fastgpt-sdk/storage build
pnpm --filter @fastgpt-sdk/storage test:unit
pnpm --filter @fastgpt-sdk/storage typecheck:test
```
发布前会执行 `prepublishOnly` 自动构建产物到 `dist/`
真实对象存储的统一契约测试位于 `sdk/storage/test/integration`。复制
`sdk/storage/.env.test.example``sdk/storage/.env.test.local`,填写凭证并将对应
`STORAGE_TEST_<PROVIDER>_ENABLED` 设置为 `true`,并配置对应的
`STORAGE_TEST_<PROVIDER>_BUCKET` 后运行。测试桶名必须以 `fastgpt-sdk-` 开头:
```bash
pnpm --filter @fastgpt-sdk/storage test:integration
pnpm --filter @fastgpt-sdk/storage test:integration:common
pnpm --filter @fastgpt-sdk/storage test:integration:minio
```
集成测试分为两层:
- `test/integration/common`:18 个 `IStorage` 通用契约,每个启用的 provider 都运行完全相同的用例。
- `test/integration/minio`:11 个 MinIO 专项用例,覆盖中断运行后的桶重建、400/1000 条分页边界、URL 编码、公共策略、真实 HTTP socket 超时,以及等待响应头和读取响应体时的下载取消。
- `test/integration/transport`:2 个无需云凭证的 OSS/COS 真实 socket 取消用例。
每个 provider 使用配置中的固定专用测试桶。每次 suite 启动时,harness 会先清空并删除可能由上次失败运行遗留的同名桶,再重新创建;结束时也会清理。不要对同一组测试配置并发运行集成测试。未启用的 provider 会被跳过。
发布前会执行 `prepublishOnly` 自动构建产物到 `dist/`
......@@ -48,6 +48,12 @@
"scripts": {
"build": "tsdown",
"dev": "tsdown --watch",
"test": "vitest run --config vitest.config.ts",
"test:unit": "vitest run --config vitest.config.ts test/unit",
"test:integration": "vitest run --config vitest.config.ts test/integration",
"test:integration:common": "vitest run --config vitest.config.ts test/integration/common",
"test:integration:minio": "vitest run --config vitest.config.ts test/integration/minio",
"typecheck:test": "tsc --noEmit -p tsconfig.test.json",
"prepublishOnly": "pnpm build"
},
"dependencies": {
......
......@@ -45,6 +45,17 @@ import type { Readable } from 'node:stream';
import { camelCase, chunk, isNotNil, kebabCase, trim } from 'es-toolkit';
import { getSignedUrl } from '@aws-sdk/s3-request-presigner';
import { DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS } from '../constants';
import {
bindAbortSignalToReadable,
encodeObjectKeyPath,
throwIfStorageDownloadAborted
} from '../utils';
import {
assertStorageObjectKey,
assertStorageObjectKeys,
assertStorageObjectPrefix,
assertRequiredStorageObjectPrefix
} from '../assert';
export class AwsS3StorageAdapter implements IStorage {
protected readonly client: S3Client;
......@@ -69,6 +80,7 @@ export class AwsS3StorageAdapter implements IStorage {
async checkObjectExists(params: ExistsObjectParams): Promise<ExistsObjectResult> {
const { key } = params;
assertStorageObjectKey(key);
let exists = false;
......@@ -97,6 +109,7 @@ export class AwsS3StorageAdapter implements IStorage {
async getObjectMetadata(params: GetObjectMetadataParams): Promise<GetObjectMetadataResult> {
const { key } = params;
assertStorageObjectKey(key);
const result = await this.client.send(
new HeadObjectCommand({
......@@ -135,6 +148,7 @@ export class AwsS3StorageAdapter implements IStorage {
async uploadObject(params: UploadObjectParams): Promise<UploadObjectResult> {
const { key, body, contentType, contentLength, contentDisposition, metadata } = params;
assertStorageObjectKey(key);
const meta: StorageObjectMetadata = {};
if (metadata) {
......@@ -167,6 +181,8 @@ export class AwsS3StorageAdapter implements IStorage {
async downloadObject(params: DownloadObjectParams): Promise<DownloadObjectResult> {
const { key, abortSignal } = params;
assertStorageObjectKey(key);
throwIfStorageDownloadAborted(abortSignal);
const result = await this.client.send(
new GetObjectCommand({
......@@ -179,16 +195,19 @@ export class AwsS3StorageAdapter implements IStorage {
if (!result.Body) {
throw new EmptyObjectError('Object is undefined');
}
const body = result.Body as Readable;
bindAbortSignalToReadable({ readable: body, abortSignal });
return {
key,
bucket: this.options.bucket,
body: result.Body as Readable
body
};
}
async deleteObject(params: DeleteObjectParams): Promise<DeleteObjectResult> {
const { key } = params;
assertStorageObjectKey(key);
await this.client.send(
new DeleteObjectCommand({
......@@ -205,6 +224,7 @@ export class AwsS3StorageAdapter implements IStorage {
async deleteObjectsByMultiKeys(params: DeleteObjectsParams): Promise<DeleteObjectsResult> {
const { keys } = params;
assertStorageObjectKeys(keys);
if (keys.length === 0) {
return {
......@@ -237,9 +257,7 @@ export class AwsS3StorageAdapter implements IStorage {
async deleteObjectsByPrefix(params: DeleteObjectsByPrefixParams): Promise<DeleteObjectsResult> {
const { prefix } = params;
if (!prefix) {
throw new Error('Prefix is required');
}
assertRequiredStorageObjectPrefix(prefix);
const fails: StorageObjectKey[] = [];
let isTruncated = false;
......@@ -258,7 +276,7 @@ export class AwsS3StorageAdapter implements IStorage {
if (!listResponse.Contents || listResponse.Contents.length === 0) {
return {
bucket: this.options.bucket,
keys: []
keys: fails
};
}
......@@ -287,6 +305,7 @@ export class AwsS3StorageAdapter implements IStorage {
async generatePresignedPutUrl(params: PresignedPutUrlParams): Promise<PresignedPutUrlResult> {
const { key, expiredSeconds, metadata, contentType } = params;
assertStorageObjectKey(key);
const expiresIn = expiredSeconds ? expiredSeconds : DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS;
......@@ -341,6 +360,7 @@ export class AwsS3StorageAdapter implements IStorage {
async generatePresignedGetUrl(params: PresignedGetUrlParams): Promise<PresignedGetUrlResult> {
const { key, expiredSeconds, responseContentType } = params;
assertStorageObjectKey(key);
const expiresIn = expiredSeconds ? expiredSeconds : DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS;
......@@ -365,13 +385,15 @@ export class AwsS3StorageAdapter implements IStorage {
generatePublicGetUrl(params: GeneratePublicGetUrlParams): GeneratePublicGetUrlResult {
const { key } = params;
assertStorageObjectKey(key);
const encodedKey = encodeObjectKeyPath(key);
let url: string;
if (this.options.forcePathStyle) {
if (this.options.publicAccessExtraSubPath) {
url = `${this.options.endpoint}/${trim(this.options.publicAccessExtraSubPath, '/')}/${this.options.bucket}/${key}`;
url = `${this.options.endpoint}/${trim(this.options.publicAccessExtraSubPath, '/')}/${this.options.bucket}/${encodedKey}`;
} else {
url = `${this.options.endpoint}/${this.options.bucket}/${key}`;
url = `${this.options.endpoint}/${this.options.bucket}/${encodedKey}`;
}
} else {
const endpoint = new URL(this.options.endpoint);
......@@ -379,9 +401,9 @@ export class AwsS3StorageAdapter implements IStorage {
const host = endpoint.host;
if (this.options.publicAccessExtraSubPath) {
url = `${protocol}//${this.options.bucket}.${host}/${trim(this.options.publicAccessExtraSubPath, '/')}/${key}`;
url = `${protocol}//${this.options.bucket}.${host}/${trim(this.options.publicAccessExtraSubPath, '/')}/${encodedKey}`;
} else {
url = `${protocol}//${this.options.bucket}.${host}/${key}`;
url = `${protocol}//${this.options.bucket}.${host}/${encodedKey}`;
}
}
......@@ -394,6 +416,7 @@ export class AwsS3StorageAdapter implements IStorage {
async listObjects(params: ListObjectsParams): Promise<ListObjectsResult> {
const { prefix } = params;
assertStorageObjectPrefix(prefix);
let keys: StorageObjectKey[] = [];
let isTruncated = false;
......@@ -430,6 +453,8 @@ export class AwsS3StorageAdapter implements IStorage {
async copyObjectInSelfBucket(params: CopyObjectParams): Promise<CopyObjectResult> {
const { sourceKey, targetKey } = params;
assertStorageObjectKey(sourceKey, 'sourceKey');
assertStorageObjectKey(targetKey, 'targetKey');
const encodedSourceKey = sourceKey
.split('/')
......
......@@ -31,6 +31,17 @@ import type {
import { PassThrough } from 'node:stream';
import { camelCase, isError, isNotNil, kebabCase } from 'es-toolkit';
import { DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS } from '../constants';
import {
bindAbortSignalToReadable,
encodeObjectKeyPath,
throwIfStorageDownloadAborted
} from '../utils';
import {
assertStorageObjectKey,
assertStorageObjectKeys,
assertStorageObjectPrefix,
assertRequiredStorageObjectPrefix
} from '../assert';
export class CosStorageAdapter implements IStorage {
protected readonly client: COS;
......@@ -62,6 +73,7 @@ export class CosStorageAdapter implements IStorage {
async checkObjectExists(params: ExistsObjectParams): Promise<ExistsObjectResult> {
const { key } = params;
assertStorageObjectKey(key);
let exists = false;
await new Promise<void>((resolve, reject) => {
......@@ -96,6 +108,7 @@ export class CosStorageAdapter implements IStorage {
async getObjectMetadata(params: GetObjectMetadataParams): Promise<GetObjectMetadataResult> {
const { key } = params;
assertStorageObjectKey(key);
const result = await new Promise<COS.HeadObjectResult>((resolve, reject) => {
this.client.headObject(
......@@ -161,6 +174,7 @@ export class CosStorageAdapter implements IStorage {
async uploadObject(params: UploadObjectParams): Promise<UploadObjectResult> {
const { key, body, contentType, contentLength, contentDisposition, metadata } = params;
assertStorageObjectKey(key);
const headers: Record<string, string> = {};
if (contentDisposition) headers['Content-Disposition'] = contentDisposition;
......@@ -199,17 +213,11 @@ export class CosStorageAdapter implements IStorage {
}
async downloadObject(params: DownloadObjectParams): Promise<DownloadObjectResult> {
params.abortSignal?.throwIfAborted();
assertStorageObjectKey(params.key);
throwIfStorageDownloadAborted(params.abortSignal);
const passThrough = new PassThrough();
const abortDownload = () => {
passThrough.destroy();
};
params.abortSignal?.addEventListener('abort', abortDownload, { once: true });
passThrough.once('close', () => {
params.abortSignal?.removeEventListener('abort', abortDownload);
});
bindAbortSignalToReadable({ readable: passThrough, abortSignal: params.abortSignal });
this.client.getObject(
{
......@@ -234,6 +242,7 @@ export class CosStorageAdapter implements IStorage {
async deleteObject(params: DeleteObjectParams): Promise<DeleteObjectResult> {
const { key } = params;
assertStorageObjectKey(key);
await new Promise<COS.DeleteObjectResult>((resolve, reject) => {
this.client.deleteObject(
......@@ -259,6 +268,14 @@ export class CosStorageAdapter implements IStorage {
async deleteObjectsByMultiKeys(params: DeleteObjectsParams): Promise<DeleteObjectsResult> {
const { keys } = params;
assertStorageObjectKeys(keys);
if (keys.length === 0) {
return {
bucket: this.options.bucket,
keys: []
};
}
const result = await new Promise<COS.DeleteMultipleObjectResult>((resolve, reject) => {
this.client.deleteMultipleObject(
......@@ -284,9 +301,7 @@ export class CosStorageAdapter implements IStorage {
async deleteObjectsByPrefix(params: DeleteObjectsByPrefixParams): Promise<DeleteObjectsResult> {
const { prefix } = params;
if (!prefix) {
throw new Error('Prefix is required');
}
assertRequiredStorageObjectPrefix(prefix);
const fails: StorageObjectKey[] = [];
let marker: string | undefined = undefined;
......@@ -354,6 +369,7 @@ export class CosStorageAdapter implements IStorage {
async generatePresignedPutUrl(params: PresignedPutUrlParams): Promise<PresignedPutUrlResult> {
const { key, expiredSeconds, metadata, contentType } = params;
assertStorageObjectKey(key);
const expiresIn = expiredSeconds ? expiredSeconds : DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS;
......@@ -398,6 +414,7 @@ export class CosStorageAdapter implements IStorage {
async generatePresignedGetUrl(params: PresignedGetUrlParams): Promise<PresignedGetUrlResult> {
const { key, expiredSeconds, responseContentType } = params;
assertStorageObjectKey(key);
const expiresIn = expiredSeconds ? expiredSeconds : DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS;
const url = await new Promise<string>((resolve, reject) => {
......@@ -431,12 +448,14 @@ export class CosStorageAdapter implements IStorage {
generatePublicGetUrl(params: GeneratePublicGetUrlParams): GeneratePublicGetUrlResult {
const { key } = params;
assertStorageObjectKey(key);
const encodedKey = encodeObjectKeyPath(key);
let url: string;
if (this.options.domain) {
url = `${this.options.protocol}//${this.options.domain}/${key}`;
url = `${this.options.protocol}//${this.options.domain}/${encodedKey}`;
} else {
url = `${this.options.protocol}//${this.options.bucket}.cos.${this.options.region}.myqcloud.com/${key}`;
url = `${this.options.protocol}//${this.options.bucket}.cos.${this.options.region}.myqcloud.com/${encodedKey}`;
}
return {
......@@ -448,6 +467,7 @@ export class CosStorageAdapter implements IStorage {
async listObjects(params: ListObjectsParams): Promise<ListObjectsResult> {
const { prefix } = params;
assertStorageObjectPrefix(prefix);
let keys: StorageObjectKey[] = [];
let marker: string | undefined = undefined;
......@@ -490,6 +510,8 @@ export class CosStorageAdapter implements IStorage {
async copyObjectInSelfBucket(params: CopyObjectParams): Promise<CopyObjectResult> {
const { sourceKey, targetKey } = params;
assertStorageObjectKey(sourceKey, 'sourceKey');
assertStorageObjectKey(targetKey, 'targetKey');
const encodedSourceKey = sourceKey
.split('/')
......
......@@ -31,6 +31,17 @@ import type {
import type { Readable } from 'node:stream';
import { camelCase, difference, kebabCase } from 'es-toolkit';
import { DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS } from '../constants';
import {
bindAbortSignalToReadable,
encodeObjectKeyPath,
throwIfStorageDownloadAborted
} from '../utils';
import {
assertStorageObjectKey,
assertStorageObjectKeys,
assertStorageObjectPrefix,
assertRequiredStorageObjectPrefix
} from '../assert';
export class OssStorageAdapter implements IStorage {
protected readonly client: OSS;
......@@ -61,6 +72,7 @@ export class OssStorageAdapter implements IStorage {
async checkObjectExists(params: ExistsObjectParams): Promise<ExistsObjectResult> {
const { key } = params;
assertStorageObjectKey(key);
let exists = false;
try {
......@@ -83,6 +95,7 @@ export class OssStorageAdapter implements IStorage {
async getObjectMetadata(params: GetObjectMetadataParams): Promise<GetObjectMetadataResult> {
const { key } = params;
assertStorageObjectKey(key);
const result = await this.client.head(key);
......@@ -122,6 +135,7 @@ export class OssStorageAdapter implements IStorage {
async uploadObject(params: UploadObjectParams): Promise<UploadObjectResult> {
const { key, body, contentType, contentLength, contentDisposition, metadata } = params;
assertStorageObjectKey(key);
const headers: Record<string, any> = {
'x-oss-storage-class': 'Standard',
......@@ -153,24 +167,12 @@ export class OssStorageAdapter implements IStorage {
async downloadObject(params: DownloadObjectParams): Promise<DownloadObjectResult> {
const { key, abortSignal } = params;
abortSignal?.throwIfAborted();
assertStorageObjectKey(key);
throwIfStorageDownloadAborted(abortSignal);
const result = await this.client.getStream(key);
const stream = result.stream as Readable;
const abortDownload = () => {
stream.destroy();
};
if (abortSignal?.aborted) {
abortDownload();
abortSignal.throwIfAborted();
}
abortSignal?.addEventListener('abort', abortDownload, { once: true });
stream.once('close', () => {
abortSignal?.removeEventListener('abort', abortDownload);
});
bindAbortSignalToReadable({ readable: stream, abortSignal });
return {
key,
......@@ -181,6 +183,7 @@ export class OssStorageAdapter implements IStorage {
async deleteObject(params: DeleteObjectParams): Promise<DeleteObjectResult> {
const { key } = params;
assertStorageObjectKey(key);
await this.client.delete(key);
......@@ -192,20 +195,46 @@ export class OssStorageAdapter implements IStorage {
async deleteObjectsByMultiKeys(params: DeleteObjectsParams): Promise<DeleteObjectsResult> {
const { keys } = params;
assertStorageObjectKeys(keys);
const result = await this.client.deleteMulti(keys, { quiet: true });
if (keys.length === 0) {
return {
bucket: this.options.bucket,
keys: []
};
}
// verbose 模式会返回成功删除的 key;quiet 全成功时响应为空,无法与失败区分。
const result = await this.client.deleteMulti(keys, { quiet: false });
const deletedKeys = (() => {
const deletedItems: unknown = result.deleted;
if (!Array.isArray(deletedItems)) return [];
const normalizedKeys: string[] = [];
for (const item of deletedItems) {
if (typeof item === 'string') {
normalizedKeys.push(item);
continue;
}
// ali-oss 的类型声明是 string[],但标准 OSS XML 在运行时解析为 { Key }[]。
if (item && typeof item === 'object' && 'Key' in item && typeof item.Key === 'string') {
normalizedKeys.push(item.Key);
continue;
}
return [];
}
return normalizedKeys;
})();
return {
bucket: this.options.bucket,
keys: difference(keys, result.deleted ?? [])
keys: difference(keys, deletedKeys)
};
}
async deleteObjectsByPrefix(params: DeleteObjectsByPrefixParams): Promise<DeleteObjectsResult> {
const { prefix } = params;
if (!prefix) {
throw new Error('Prefix is required');
}
assertRequiredStorageObjectPrefix(prefix);
const fails: StorageObjectKey[] = [];
let marker: string | undefined = undefined;
......@@ -226,7 +255,7 @@ export class OssStorageAdapter implements IStorage {
if (!listResponse.objects || listResponse.objects.length === 0) {
return {
bucket: this.options.bucket,
keys: []
keys: fails
};
}
......@@ -247,6 +276,7 @@ export class OssStorageAdapter implements IStorage {
async generatePresignedPutUrl(params: PresignedPutUrlParams): Promise<PresignedPutUrlResult> {
const { key, expiredSeconds, metadata, contentType } = params;
assertStorageObjectKey(key);
const expiresIn = expiredSeconds ? expiredSeconds : DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS;
......@@ -285,6 +315,7 @@ export class OssStorageAdapter implements IStorage {
async generatePresignedGetUrl(params: PresignedGetUrlParams): Promise<PresignedGetUrlResult> {
const { key, expiredSeconds, responseContentType } = params;
assertStorageObjectKey(key);
const expiresIn = expiredSeconds ? expiredSeconds : DEFAULT_PRESIGNED_URL_EXPIRED_SECONDS;
const url = this.client.signatureUrl(key, {
......@@ -308,6 +339,8 @@ export class OssStorageAdapter implements IStorage {
generatePublicGetUrl(params: GeneratePublicGetUrlParams): GeneratePublicGetUrlResult {
const { key } = params;
assertStorageObjectKey(key);
const encodedKey = encodeObjectKeyPath(key);
let protocol = 'https:';
if (!this.options.secure) {
......@@ -316,9 +349,9 @@ export class OssStorageAdapter implements IStorage {
let url: string;
if (this.options.cname) {
url = `${protocol}//${this.options.endpoint}/${key}`;
url = `${protocol}//${this.options.endpoint}/${encodedKey}`;
} else {
url = `${protocol}//${this.options.bucket}.${this.options.region}.aliyuncs.com/${key}`;
url = `${protocol}//${this.options.bucket}.${this.options.region}.aliyuncs.com/${encodedKey}`;
}
return {
......@@ -330,6 +363,7 @@ export class OssStorageAdapter implements IStorage {
async listObjects(params: ListObjectsParams): Promise<ListObjectsResult> {
const { prefix } = params;
assertStorageObjectPrefix(prefix);
let keys: StorageObjectKey[] = [];
let marker: string | undefined = undefined;
......@@ -367,6 +401,8 @@ export class OssStorageAdapter implements IStorage {
async copyObjectInSelfBucket(params: CopyObjectParams): Promise<CopyObjectResult> {
const { sourceKey, targetKey } = params;
assertStorageObjectKey(sourceKey, 'sourceKey');
assertStorageObjectKey(targetKey, 'targetKey');
await this.client.copy(targetKey, sourceKey);
......
import { InvalidStorageObjectKeyError, type InvalidStorageObjectKeyReason } from './errors';
/** 四个 adapter 都可移植的对象 key 最大 UTF-8 字节数。 */
export const MAX_STORAGE_OBJECT_KEY_UTF8_BYTES = 850;
function throwInvalidStorageObjectKey({
field,
reason,
actualBytes
}: {
field: string;
reason: InvalidStorageObjectKeyReason;
actualBytes?: number;
}): never {
throw new InvalidStorageObjectKeyError({
field,
reason,
actualBytes,
maxBytes: actualBytes === undefined ? undefined : MAX_STORAGE_OBJECT_KEY_UTF8_BYTES
});
}
/**
* 检查字符串是否不存在未配对的 UTF-16 surrogate。
* Buffer 会把非法 surrogate 静默替换成 U+FFFD,因此必须在计算 UTF-8 长度前显式检查。
*/
function isWellFormedUnicode(value: string): boolean {
for (let index = 0; index < value.length; index += 1) {
const codeUnit = value.charCodeAt(index);
if (codeUnit >= 0xd800 && codeUnit <= 0xdbff) {
const nextCodeUnit = value.charCodeAt(index + 1);
if (nextCodeUnit < 0xdc00 || nextCodeUnit > 0xdfff) return false;
index += 1;
continue;
}
if (codeUnit >= 0xdc00 && codeUnit <= 0xdfff) return false;
}
return true;
}
/**
* 按 SDK 统一规范预检对象 key;失败时不会把原始 key 写入错误消息。
* 该规范取 AWS S3、MinIO、OSS、COS 可稳定处理范围的交集。
*/
export function assertStorageObjectKey(value: unknown, field = 'key'): asserts value is string {
if (typeof value !== 'string') {
throwInvalidStorageObjectKey({ field, reason: 'invalid_type' });
}
if (value.length === 0) {
throwInvalidStorageObjectKey({ field, reason: 'empty' });
}
if (!isWellFormedUnicode(value)) {
throwInvalidStorageObjectKey({ field, reason: 'invalid_unicode' });
}
const actualBytes = Buffer.byteLength(value, 'utf8');
if (actualBytes > MAX_STORAGE_OBJECT_KEY_UTF8_BYTES) {
throwInvalidStorageObjectKey({ field, reason: 'too_long', actualBytes });
}
if (value.startsWith('/')) {
throwInvalidStorageObjectKey({ field, reason: 'leading_slash' });
}
if (value.includes('\\')) {
throwInvalidStorageObjectKey({ field, reason: 'backslash' });
}
if (value.includes('//')) {
throwInvalidStorageObjectKey({ field, reason: 'empty_path_segment' });
}
if (/[\u0000-\u001f\u007f]/u.test(value)) {
throwInvalidStorageObjectKey({ field, reason: 'control_character' });
}
if (
value.split('/').some((segment) => {
const trimmedSegment = segment.trim();
return trimmedSegment === '.' || trimmedSegment === '..';
})
) {
throwInvalidStorageObjectKey({ field, reason: 'dot_path_segment' });
}
}
/** 批量方法必须完整预检数组后,调用方才能开始分块或产生远端副作用。 */
export function assertStorageObjectKeys(keys: unknown): asserts keys is string[] {
if (!Array.isArray(keys)) {
throwInvalidStorageObjectKey({ field: 'keys', reason: 'invalid_type' });
}
for (let index = 0; index < keys.length; index += 1) {
assertStorageObjectKey(keys[index], `keys[${index}]`);
}
}
/** listObjects 允许省略或传空 prefix;非空 prefix 与对象 key 使用同一规范。 */
export function assertStorageObjectPrefix(prefix: unknown): asserts prefix is string | undefined {
if (prefix === undefined || prefix === '') return;
assertStorageObjectKey(prefix, 'prefix');
}
/** 删除前缀必须非空;其余字符和长度限制与对象 key 完全一致。 */
export function assertRequiredStorageObjectPrefix(prefix: unknown): asserts prefix is string {
if (typeof prefix === 'string' && prefix.trim().length === 0) {
throw new Error('Prefix is required');
}
assertStorageObjectKey(prefix, 'prefix');
}
......@@ -18,3 +18,59 @@ export class EmptyObjectError extends Error {
this.name = 'EmptyObjectError';
}
}
export type InvalidStorageObjectKeyReason =
| 'invalid_type'
| 'empty'
| 'invalid_unicode'
| 'too_long'
| 'leading_slash'
| 'backslash'
| 'empty_path_segment'
| 'dot_path_segment'
| 'control_character';
const invalidStorageObjectKeyReasonMessages: Record<InvalidStorageObjectKeyReason, string> = {
invalid_type: 'must be a string',
empty: 'must not be empty',
invalid_unicode: 'must contain well-formed Unicode',
too_long: 'exceeds the UTF-8 byte limit',
leading_slash: 'must not start with a slash',
backslash: 'must not contain a backslash',
empty_path_segment: 'must not contain consecutive slashes',
dot_path_segment: 'must not contain dot path segments',
control_character: 'must not contain ASCII control characters'
};
/** SDK 在远端请求前发现对象 key 或 prefix 不符合统一可移植规范。 */
export class InvalidStorageObjectKeyError extends Error {
readonly field: string;
readonly reason: InvalidStorageObjectKeyReason;
readonly actualBytes?: number;
readonly maxBytes?: number;
constructor({
field,
reason,
actualBytes,
maxBytes
}: {
field: string;
reason: InvalidStorageObjectKeyReason;
actualBytes?: number;
maxBytes?: number;
}) {
const byteDetails =
actualBytes !== undefined && maxBytes !== undefined
? ` (${actualBytes} bytes, maximum ${maxBytes})`
: '';
super(
`Invalid storage object ${field}: ${invalidStorageObjectKeyReasonMessages[reason]}${byteDetails}`
);
this.name = 'InvalidStorageObjectKeyError';
this.field = field;
this.reason = reason;
this.actualBytes = actualBytes;
this.maxBytes = maxBytes;
}
}
import type { Readable } from 'node:stream';
import { Readable as NodeReadable } from 'node:stream';
import type { Mock } from 'vitest';
import type { IStorage } from '../interface';
import type {
CopyObjectParams,
CopyObjectResult,
DeleteObjectParams,
DeleteObjectResult,
DeleteObjectsByPrefixParams,
DeleteObjectsParams,
DeleteObjectsResult,
DownloadObjectParams,
DownloadObjectResult,
EnsureBucketResult,
ExistsObjectParams,
ExistsObjectResult,
GeneratePublicGetUrlParams,
GeneratePublicGetUrlResult,
GetObjectMetadataParams,
GetObjectMetadataResult,
ListObjectsParams,
ListObjectsResult,
PresignedGetUrlParams,
PresignedGetUrlResult,
PresignedPutUrlParams,
PresignedPutUrlResult,
StorageObjectKey,
StorageObjectMetadata,
StorageUploadBody,
UploadObjectParams,
UploadObjectResult
} from '../types';
import {
assertStorageObjectKey,
assertStorageObjectKeys,
assertStorageObjectPrefix,
assertRequiredStorageObjectPrefix
} from '../assert';
import { bindAbortSignalToReadable, throwIfStorageDownloadAborted } from '../utils';
type VitestLike = {
fn: <T extends (...args: any[]) => any>(impl?: T) => Mock<T>;
};
type StoredObject = {
body: Buffer;
metadata: StorageObjectMetadata;
contentType?: string;
contentLength?: number;
contentDisposition?: string;
etag?: string;
};
export type VitestStorageMock = IStorage & {
/** 便于在测试中直接读写内存对象(key -> object)。 */
__objects: Map<StorageObjectKey, StoredObject>;
/** 清空内存对象。 */
__reset: () => void;
/** 直接写入一个对象(绕过 uploadObject)。 */
__putObject: (key: StorageObjectKey, obj: Partial<StoredObject> & { body: Buffer }) => void;
};
export type CreateVitestStorageMockParams = {
vi: VitestLike;
bucketName?: string;
/**
* 用于构造 presigned/public URL 的 base(仅 mock 用)。
* 例如:`https://mock-storage.local`
*/
baseUrl?: string;
};
async function bodyToBuffer(body: StorageUploadBody): Promise<Buffer> {
if (Buffer.isBuffer(body)) return body;
if (typeof body === 'string') return Buffer.from(body);
return await readableToBuffer(body);
}
async function readableToBuffer(readable: Readable): Promise<Buffer> {
const chunks: Buffer[] = [];
for await (const chunk of readable) {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
}
return Buffer.concat(chunks);
}
function bufferToReadable(buf: Buffer): Readable {
return NodeReadable.from(buf);
}
function getEtag(buf: Buffer) {
// mock: 非加密 hash,只是为了在测试里有稳定值可断言
return `etag_${buf.length}_${buf.subarray(0, 8).toString('hex')}`;
}
export function createVitestStorageMock(params: CreateVitestStorageMockParams): VitestStorageMock {
const { vi, bucketName = 'mock-bucket', baseUrl = 'https://mock-storage.local' } = params;
const objects = new Map<StorageObjectKey, StoredObject>();
let bucketEnsured = false;
const ensureBucket = vi.fn(async (): Promise<EnsureBucketResult> => {
const exists = bucketEnsured;
bucketEnsured = true;
return { exists, created: !exists, bucket: bucketName };
});
const checkObjectExists = vi.fn(
async ({ key }: ExistsObjectParams): Promise<ExistsObjectResult> => {
assertStorageObjectKey(key);
return { bucket: bucketName, key, exists: objects.has(key) };
}
);
const uploadObject = vi.fn(async (p: UploadObjectParams): Promise<UploadObjectResult> => {
assertStorageObjectKey(p.key);
const buf = await bodyToBuffer(p.body);
const contentLength = p.contentLength ?? buf.length;
objects.set(p.key, {
body: buf,
metadata: p.metadata ?? {},
contentType: p.contentType,
contentDisposition: p.contentDisposition,
contentLength,
etag: getEtag(buf)
});
return { bucket: bucketName, key: p.key };
});
const downloadObject = vi.fn(async (p: DownloadObjectParams): Promise<DownloadObjectResult> => {
assertStorageObjectKey(p.key);
throwIfStorageDownloadAborted(p.abortSignal);
const obj = objects.get(p.key);
if (!obj) {
throw new Error(`Object not found: ${p.key}`);
}
const body = bufferToReadable(obj.body);
bindAbortSignalToReadable({ readable: body, abortSignal: p.abortSignal });
return { bucket: bucketName, key: p.key, body };
});
const deleteObject = vi.fn(async (p: DeleteObjectParams): Promise<DeleteObjectResult> => {
assertStorageObjectKey(p.key);
objects.delete(p.key);
return { bucket: bucketName, key: p.key };
});
const deleteObjectsByMultiKeys = vi.fn(
async (p: DeleteObjectsParams): Promise<DeleteObjectsResult> => {
assertStorageObjectKeys(p.keys);
for (const key of p.keys) objects.delete(key);
return { bucket: bucketName, keys: [] };
}
);
const deleteObjectsByPrefix = vi.fn(
async (p: DeleteObjectsByPrefixParams): Promise<DeleteObjectsResult> => {
assertRequiredStorageObjectPrefix(p.prefix);
const keys: string[] = [];
for (const key of objects.keys()) {
if (key.startsWith(p.prefix)) keys.push(key);
}
for (const key of keys) objects.delete(key);
return { bucket: bucketName, keys: [] };
}
);
const generatePresignedPutUrl = vi.fn(
async (p: PresignedPutUrlParams): Promise<PresignedPutUrlResult> => {
assertStorageObjectKey(p.key);
const putUrl = `${baseUrl}/put/${encodeURIComponent(bucketName)}/${encodeURIComponent(p.key)}`;
// mock: 直接透传 metadata 作为“headers”
const metadata: Record<string, string> = p.metadata ? { ...p.metadata } : {};
return { bucket: bucketName, key: p.key, url: putUrl, metadata };
}
);
const generatePresignedGetUrl = vi.fn(
async (p: PresignedGetUrlParams): Promise<PresignedGetUrlResult> => {
assertStorageObjectKey(p.key);
const query = p.responseContentType
? `?response-content-type=${encodeURIComponent(p.responseContentType)}`
: '';
const getUrl = `${baseUrl}/get/${encodeURIComponent(bucketName)}/${encodeURIComponent(p.key)}${query}`;
return { bucket: bucketName, key: p.key, url: getUrl };
}
);
const generatePublicGetUrl = vi.fn(
({ key }: GeneratePublicGetUrlParams): GeneratePublicGetUrlResult => {
assertStorageObjectKey(key);
const publicGetUrl = `${baseUrl}/public/${encodeURIComponent(bucketName)}/${encodeURIComponent(key)}`;
return { url: publicGetUrl, bucket: bucketName, key };
}
);
const listObjects = vi.fn(async (p: ListObjectsParams): Promise<ListObjectsResult> => {
assertStorageObjectPrefix(p.prefix);
const keys = Array.from(objects.keys()).filter((k) =>
p.prefix ? k.startsWith(p.prefix) : true
);
keys.sort();
return { bucket: bucketName, keys };
});
const copyObjectInSelfBucket = vi.fn(async (p: CopyObjectParams): Promise<CopyObjectResult> => {
assertStorageObjectKey(p.sourceKey, 'sourceKey');
assertStorageObjectKey(p.targetKey, 'targetKey');
const src = objects.get(p.sourceKey);
if (!src) {
throw new Error(`Source object not found: ${p.sourceKey}`);
}
objects.set(p.targetKey, { ...src, body: Buffer.from(src.body) });
return { bucket: bucketName, sourceKey: p.sourceKey, targetKey: p.targetKey };
});
const getObjectMetadata = vi.fn(
async (p: GetObjectMetadataParams): Promise<GetObjectMetadataResult> => {
assertStorageObjectKey(p.key);
const obj = objects.get(p.key);
if (!obj) {
throw new Error(`Object not found: ${p.key}`);
}
return {
bucket: bucketName,
key: p.key,
metadata: obj.metadata ?? {},
contentType: obj.contentType,
contentLength: obj.contentLength,
etag: obj.etag
};
}
);
const destroy = vi.fn(async (): Promise<void> => {});
const mock: VitestStorageMock = {
bucketName,
ensureBucket,
checkObjectExists,
uploadObject,
downloadObject,
deleteObject,
deleteObjectsByMultiKeys,
deleteObjectsByPrefix,
generatePresignedPutUrl,
generatePresignedGetUrl,
generatePublicGetUrl,
listObjects,
copyObjectInSelfBucket,
getObjectMetadata,
destroy,
__objects: objects,
__reset: () => objects.clear(),
__putObject: (key, obj) => {
objects.set(key, {
body: obj.body,
metadata: obj.metadata ?? {},
contentType: obj.contentType,
contentLength: obj.contentLength ?? obj.body.length,
contentDisposition: obj.contentDisposition,
etag: obj.etag ?? getEtag(obj.body)
});
}
};
return mock;
}
export { createStorage } from './factory';
export { createVitestStorageMock } from './testing/vitestMock';
export type { VitestStorageMock, CreateVitestStorageMockParams } from './testing/vitestMock';
export { createVitestStorageMock } from './helper/mock';
export type { VitestStorageMock, CreateVitestStorageMockParams } from './helper/mock';
export type {
IStorage,
IStorageOptions,
......@@ -35,7 +35,19 @@ export type {
GetObjectMetadataParams,
GetObjectMetadataResult
} from './types';
export { NoSuchBucketError, NoBucketReadPermissionError, EmptyObjectError } from './errors';
export {
NoSuchBucketError,
NoBucketReadPermissionError,
EmptyObjectError,
InvalidStorageObjectKeyError
} from './errors';
export type { InvalidStorageObjectKeyReason } from './errors';
export {
MAX_STORAGE_OBJECT_KEY_UTF8_BYTES,
assertStorageObjectKey,
assertStorageObjectKeys,
assertStorageObjectPrefix
} from './assert';
export { AwsS3StorageAdapter } from './adapters/aws-s3.adapter';
export { CosStorageAdapter } from './adapters/cos.adapter';
export { MinioStorageAdapter } from './adapters/minio.adapter';
......
......@@ -309,7 +309,7 @@ export interface IStorage {
*
* 注意:
* - 各厂商对单次批量删除的最大数量限制不同,adapter 可能需要分批处理。
* - 返回的 `deleted` 通常只包含实际删除/确认删除的 key
* - 返回的 `keys` 只包含删除失败、需要上层重试的 key;空数组表示全部成功
*/
deleteObjectsByMultiKeys(params: DeleteObjectsParams): Promise<DeleteObjectsResult>;
......
import type { Readable } from 'node:stream';
import { Readable as NodeReadable } from 'node:stream';
import type { MockedFunction } from 'vitest';
import type { IStorage } from '../interface';
import type {
CopyObjectParams,
CopyObjectResult,
DeleteObjectParams,
DeleteObjectResult,
DeleteObjectsByPrefixParams,
DeleteObjectsParams,
DeleteObjectsResult,
DownloadObjectParams,
DownloadObjectResult,
EnsureBucketResult,
ExistsObjectParams,
ExistsObjectResult,
GeneratePublicGetUrlParams,
GeneratePublicGetUrlResult,
GetObjectMetadataParams,
GetObjectMetadataResult,
ListObjectsParams,
ListObjectsResult,
PresignedGetUrlParams,
PresignedGetUrlResult,
PresignedPutUrlParams,
PresignedPutUrlResult,
StorageObjectKey,
StorageObjectMetadata,
StorageUploadBody,
UploadObjectParams,
UploadObjectResult
} from '../types';
type VitestLike = {
fn: <T extends (...args: any[]) => any>(impl?: T) => MockedFunction<T>;
};
type StoredObject = {
body: Buffer;
metadata: StorageObjectMetadata;
contentType?: string;
contentLength?: number;
contentDisposition?: string;
etag?: string;
};
export type VitestStorageMock = IStorage & {
/** 便于在测试中直接读写内存对象(key -> object)。 */
__objects: Map<StorageObjectKey, StoredObject>;
/** 清空内存对象。 */
__reset: () => void;
/** 直接写入一个对象(绕过 uploadObject)。 */
__putObject: (key: StorageObjectKey, obj: Partial<StoredObject> & { body: Buffer }) => void;
};
export type CreateVitestStorageMockParams = {
vi: VitestLike;
bucketName?: string;
/**
* 用于构造 presigned/public URL 的 base(仅 mock 用)。
* 例如:`https://mock-storage.local`
*/
baseUrl?: string;
};
async function bodyToBuffer(body: StorageUploadBody): Promise<Buffer> {
if (Buffer.isBuffer(body)) return body;
if (typeof body === 'string') return Buffer.from(body);
return await readableToBuffer(body);
}
async function readableToBuffer(readable: Readable): Promise<Buffer> {
const chunks: Buffer[] = [];
for await (const chunk of readable) {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
}
return Buffer.concat(chunks);
}
function bufferToReadable(buf: Buffer): Readable {
return NodeReadable.from(buf);
}
function getEtag(buf: Buffer) {
// mock: 非加密 hash,只是为了在测试里有稳定值可断言
return `etag_${buf.length}_${buf.subarray(0, 8).toString('hex')}`;
}
export function createVitestStorageMock(params: CreateVitestStorageMockParams): VitestStorageMock {
const { vi, bucketName = 'mock-bucket', baseUrl = 'https://mock-storage.local' } = params;
const objects = new Map<StorageObjectKey, StoredObject>();
let bucketEnsured = false;
const ensureBucket = vi.fn(async (): Promise<EnsureBucketResult> => {
const exists = bucketEnsured;
bucketEnsured = true;
return { exists, created: !exists, bucket: bucketName };
});
const checkObjectExists = vi.fn(
async ({ key }: ExistsObjectParams): Promise<ExistsObjectResult> => {
return { bucket: bucketName, key, exists: objects.has(key) };
}
);
const uploadObject = vi.fn(async (p: UploadObjectParams): Promise<UploadObjectResult> => {
const buf = await bodyToBuffer(p.body);
const contentLength = p.contentLength ?? buf.length;
objects.set(p.key, {
body: buf,
metadata: p.metadata ?? {},
contentType: p.contentType,
contentDisposition: p.contentDisposition,
contentLength,
etag: getEtag(buf)
});
return { bucket: bucketName, key: p.key };
});
const downloadObject = vi.fn(async (p: DownloadObjectParams): Promise<DownloadObjectResult> => {
p.abortSignal?.throwIfAborted();
const obj = objects.get(p.key);
if (!obj) {
throw new Error(`Object not found: ${p.key}`);
}
const body = bufferToReadable(obj.body);
const abortDownload = () => {
body.destroy();
};
p.abortSignal?.addEventListener('abort', abortDownload, { once: true });
body.once('close', () => {
p.abortSignal?.removeEventListener('abort', abortDownload);
});
return { bucket: bucketName, key: p.key, body };
});
const deleteObject = vi.fn(async (p: DeleteObjectParams): Promise<DeleteObjectResult> => {
objects.delete(p.key);
return { bucket: bucketName, key: p.key };
});
const deleteObjectsByMultiKeys = vi.fn(
async (p: DeleteObjectsParams): Promise<DeleteObjectsResult> => {
for (const key of p.keys) objects.delete(key);
return { bucket: bucketName, keys: p.keys };
}
);
const deleteObjectsByPrefix = vi.fn(
async (p: DeleteObjectsByPrefixParams): Promise<DeleteObjectsResult> => {
if (!p.prefix) {
throw new Error('prefix must be a non-empty string');
}
const keys: string[] = [];
for (const key of objects.keys()) {
if (key.startsWith(p.prefix)) keys.push(key);
}
for (const key of keys) objects.delete(key);
return { bucket: bucketName, keys };
}
);
const generatePresignedPutUrl = vi.fn(
async (p: PresignedPutUrlParams): Promise<PresignedPutUrlResult> => {
const putUrl = `${baseUrl}/put/${encodeURIComponent(bucketName)}/${encodeURIComponent(p.key)}`;
// mock: 直接透传 metadata 作为“headers”
const metadata: Record<string, string> = p.metadata ? { ...p.metadata } : {};
return { bucket: bucketName, key: p.key, url: putUrl, metadata };
}
);
const generatePresignedGetUrl = vi.fn(
async (p: PresignedGetUrlParams): Promise<PresignedGetUrlResult> => {
const query = p.responseContentType
? `?response-content-type=${encodeURIComponent(p.responseContentType)}`
: '';
const getUrl = `${baseUrl}/get/${encodeURIComponent(bucketName)}/${encodeURIComponent(p.key)}${query}`;
return { bucket: bucketName, key: p.key, url: getUrl };
}
);
const generatePublicGetUrl = vi.fn(
({ key }: GeneratePublicGetUrlParams): GeneratePublicGetUrlResult => {
const publicGetUrl = `${baseUrl}/public/${encodeURIComponent(bucketName)}/${encodeURIComponent(key)}`;
return { url: publicGetUrl, bucket: bucketName, key };
}
);
const listObjects = vi.fn(async (p: ListObjectsParams): Promise<ListObjectsResult> => {
const keys = Array.from(objects.keys()).filter((k) =>
p.prefix ? k.startsWith(p.prefix) : true
);
keys.sort();
return { bucket: bucketName, keys };
});
const copyObjectInSelfBucket = vi.fn(async (p: CopyObjectParams): Promise<CopyObjectResult> => {
const src = objects.get(p.sourceKey);
if (!src) {
throw new Error(`Source object not found: ${p.sourceKey}`);
}
objects.set(p.targetKey, { ...src, body: Buffer.from(src.body) });
return { bucket: bucketName, sourceKey: p.sourceKey, targetKey: p.targetKey };
});
const getObjectMetadata = vi.fn(
async (p: GetObjectMetadataParams): Promise<GetObjectMetadataResult> => {
const obj = objects.get(p.key);
if (!obj) {
throw new Error(`Object not found: ${p.key}`);
}
return {
bucket: bucketName,
key: p.key,
metadata: obj.metadata ?? {},
contentType: obj.contentType,
contentLength: obj.contentLength,
etag: obj.etag
};
}
);
const destroy = vi.fn(async (): Promise<void> => {});
const mock: VitestStorageMock = {
bucketName,
ensureBucket,
checkObjectExists,
uploadObject,
downloadObject,
deleteObject,
deleteObjectsByMultiKeys,
deleteObjectsByPrefix,
generatePresignedPutUrl,
generatePresignedGetUrl,
generatePublicGetUrl,
listObjects,
copyObjectInSelfBucket,
getObjectMetadata,
destroy,
__objects: objects,
__reset: () => objects.clear(),
__putObject: (key, obj) => {
objects.set(key, {
body: obj.body,
metadata: obj.metadata ?? {},
contentType: obj.contentType,
contentLength: obj.contentLength ?? obj.body.length,
contentDisposition: obj.contentDisposition,
etag: obj.etag ?? getEtag(obj.body)
});
}
};
return mock;
}
// Keep the former test helper path available to external workspace consumers.
export * from '../helper/mock';
......@@ -24,6 +24,8 @@ export type StorageBucketName = string;
* 说明:
* - 在同一个 bucket 内唯一标识一个对象。
* - 通常形如:`a/b/c.txt`(用 `/` 形成“目录”层级,但对象存储并不是真正的目录结构)。
* - SDK 统一限制为 1 - 850 UTF-8 bytes,并拒绝不可移植的控制字符、反斜线、空路径段和 `.`/`..` 路径段。
* - 空格、`+`、`#`、`&`、`%`、`?`、中文和 emoji 均为合法字符,由 adapter 在 URL 层编码。
*/
export type StorageObjectKey = string;
......
import type { Readable } from 'node:stream';
/** 将对象 key 编码为 URL path,同时保留对象存储使用的 `/` 层级分隔符。 */
export const encodeObjectKeyPath = (key: string): string =>
key
.split('/')
.map((segment) => encodeURIComponent(segment))
.join('/');
/** 将任意 AbortSignal.reason 归一为可用于 Readable.destroy 的 Error。 */
export const getAbortSignalError = (abortSignal: AbortSignal): Error => {
if (abortSignal.reason instanceof Error) return abortSignal.reason;
const error = new Error(
abortSignal.reason === undefined ? 'The operation was aborted' : String(abortSignal.reason)
);
error.name = 'AbortError';
return error;
};
/** 在发起远端下载前检查取消,确保预取消请求不会进入厂商 SDK。 */
export const throwIfStorageDownloadAborted = (abortSignal?: AbortSignal): void => {
if (abortSignal?.aborted) throw getAbortSignalError(abortSignal);
};
/**
* 将下载取消绑定到返回流,并在流关闭后解除监听。
* 二次检查覆盖等待厂商返回流期间发生的取消竞态。
*/
export const bindAbortSignalToReadable = ({
readable,
abortSignal
}: {
readable: Readable;
abortSignal?: AbortSignal;
}): void => {
if (!abortSignal) return;
if (abortSignal.aborted) {
readable.destroy();
throw getAbortSignalError(abortSignal);
}
const abortDownload = () => readable.destroy(getAbortSignalError(abortSignal));
abortSignal.addEventListener('abort', abortDownload, { once: true });
readable.once('close', () => {
abortSignal.removeEventListener('abort', abortDownload);
});
};
import fs from 'node:fs';
/**
* Priority: .env.test.local > .env.test > .env.local > .env
*/
function getEnvFilePath(): URL | undefined {
const files = ['.env.test.local', '.env.test', '.env.local', '.env'];
return files.map((f) => new URL(f, import.meta.url)).find((p) => fs.existsSync(p));
}
export function setup() {
process.loadEnvFile(getEnvFilePath());
}
export function teardown() {
// no-op
}
import { storageIntegrationProviders } from '../providers';
import { runStorageAdapterContract } from './storage.contract';
for (const provider of storageIntegrationProviders) {
runStorageAdapterContract(provider);
}
import type { IStorage } from '../../src/interface';
/** 构造指定总字节数的 ASCII key,并限制单个路径段长度以兼容文件系统型对象存储。 */
export const createAsciiKeyAtLength = ({
prefix,
byteLength,
maxSegmentLength = 200
}: {
prefix: string;
byteLength: number;
maxSegmentLength?: number;
}): string => {
let remainingLength = byteLength - Buffer.byteLength(prefix);
if (remainingLength <= 0) {
throw new Error('Target byte length must be longer than the prefix');
}
const segments: string[] = [];
while (remainingLength > 0) {
const separatorLength = segments.length > 0 ? 1 : 0;
const segmentLength = Math.min(maxSegmentLength, remainingLength - separatorLength);
if (segmentLength <= 0) throw new Error('Insufficient space for another path segment');
segments.push('a'.repeat(segmentLength));
remainingLength -= segmentLength + separatorLength;
}
return `${prefix}${segments.join('/')}`;
};
/**
* 删除已存在的固定集成测试桶,供下次运行重新创建干净环境。
* `DeleteObjectsResult.keys` 是失败项;只要存在失败 key,就保留桶并让测试失败。
*/
export const removeIntegrationBucketIfExists = async ({
storage,
bucketExists,
deleteBucket
}: {
storage: IStorage;
bucketExists: () => Promise<boolean>;
deleteBucket: () => Promise<void>;
}): Promise<void> => {
if (!(await bucketExists())) return;
const { keys } = await storage.listObjects({});
if (keys.length > 0) {
const { keys: failedKeys } = await storage.deleteObjectsByMultiKeys({ keys });
if (failedKeys.length > 0) {
throw new Error(`Failed to clean integration test bucket: ${failedKeys.join(', ')}`);
}
}
await deleteBucket();
};
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import type { MinioStorageAdapter } from '../../../src/adapters/minio.adapter';
import { InvalidStorageObjectKeyError } from '../../../src/errors';
import { minioIntegrationProvider, type StorageIntegrationContext } from '../providers';
import { createAsciiKeyAtLength } from '../helpers';
const uploadInBatches = async ({
context,
keys,
batchSize
}: {
context: StorageIntegrationContext;
keys: string[];
batchSize: number;
}) => {
for (let index = 0; index < keys.length; index += batchSize) {
await Promise.all(
keys.slice(index, index + batchSize).map((key) =>
context.storage.uploadObject({
key,
body: 'x',
contentType: 'text/plain',
contentLength: 1
})
)
);
}
};
describe.skipIf(!minioIntegrationProvider.enabled).sequential('MinIO-specific integration', () => {
let context: StorageIntegrationContext;
beforeAll(async () => {
context = await minioIntegrationProvider.createContext();
});
afterAll(async () => {
await context?.cleanup();
});
it('recreates the stable bucket and removes objects left by an interrupted run', async () => {
const interruptedContext = context;
const staleKey = `${interruptedContext.rootPrefix}stale/object.txt`;
await interruptedContext.storage.uploadObject({ key: staleKey, body: 'stale' });
await interruptedContext.storage.destroy();
context = await minioIntegrationProvider.createContext();
expect(context.bucket).toBe(interruptedContext.bucket);
await expect(context.storage.listObjects({ prefix: staleKey })).resolves.toEqual({
bucket: context.bucket,
keys: []
});
});
it('creates a missing bucket through MinioStorageAdapter', () => {
expect(context.initialEnsureResult).toEqual({
bucket: context.bucket,
exists: false,
created: true
});
});
it('deletes 401 URL-sensitive keys across the 400-object prefix page boundary', async () => {
const prefix = `${context.rootPrefix}prefix-page/team & +/`;
const keys = Array.from({ length: 401 }, (_, index) => `${prefix}file + ${index}.txt`);
await uploadInBatches({ context, keys, batchSize: 20 });
const beforeDelete = await context.storage.listObjects({ prefix });
expect(new Set(beforeDelete.keys)).toEqual(new Set(keys));
await expect(context.storage.deleteObjectsByPrefix({ prefix })).resolves.toEqual({
bucket: context.bucket,
keys: []
});
await expect(context.storage.listObjects({ prefix })).resolves.toEqual({
bucket: context.bucket,
keys: []
});
});
it('lists 1001 objects across pages and deletes them across 1000-key batches', async () => {
const prefix = `${context.rootPrefix}list-page/`;
const keys = Array.from({ length: 1001 }, (_, index) => `${prefix}${index}.txt`);
await uploadInBatches({ context, keys, batchSize: 25 });
const listed = await context.storage.listObjects({ prefix });
expect(new Set(listed.keys)).toEqual(new Set(keys));
await expect(context.storage.deleteObjectsByMultiKeys({ keys })).resolves.toEqual({
bucket: context.bucket,
keys: []
});
await expect(context.storage.listObjects({ prefix })).resolves.toEqual({
bucket: context.bucket,
keys: []
});
});
it('rejects an object key beyond the portable 850-byte limit without creating an object', async () => {
const prefix = `${context.rootPrefix}too-long/`;
const key = createAsciiKeyAtLength({ prefix, byteLength: 851 });
await expect(context.storage.uploadObject({ key, body: 'too-long' })).rejects.toMatchObject({
name: InvalidStorageObjectKeyError.name,
reason: 'too_long',
actualBytes: 851,
maxBytes: 850
});
await expect(context.storage.listObjects({ prefix })).resolves.toEqual({
bucket: context.bucket,
keys: []
});
});
it('grants anonymous GET without granting anonymous PUT', async () => {
const storage = context.storage as MinioStorageAdapter;
const key = `${context.rootPrefix}public/folder name/file #+.txt`;
await storage.uploadObject({ key, body: 'public-content' });
const publicUrl = storage.generatePublicGetUrl({ key }).url;
const privateResponse = await fetch(publicUrl);
expect(privateResponse.status).toBe(403);
await storage.ensurePublicBucketPolicy();
const publicResponse = await fetch(publicUrl);
expect(publicResponse.ok).toBe(true);
await expect(publicResponse.text()).resolves.toBe('public-content');
const anonymousPut = await fetch(publicUrl, { method: 'PUT', body: 'overwritten' });
expect(anonymousPut.status).toBe(403);
const authenticatedDownload = await storage.downloadObject({ key });
const chunks: Buffer[] = [];
for await (const chunk of authenticatedDownload.body) {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
}
expect(Buffer.concat(chunks).toString()).toBe('public-content');
});
it('removes bucket lifecycle when no lifecycle configuration exists', async () => {
const storage = context.storage as MinioStorageAdapter;
await expect(storage.removeBucketLifecycle()).resolves.toBeUndefined();
});
});
import * as http from 'node:http';
import type { AddressInfo, Socket } from 'node:net';
import { afterEach, describe, expect, it } from 'vitest';
import {
createMinioTimeoutTransport,
MinioStorageAdapter
} from '../../../src/adapters/minio.adapter';
const servers = new Set<http.Server>();
const listen = async (server: http.Server): Promise<number> => {
await new Promise<void>((resolve, reject) => {
server.once('error', reject);
server.listen(0, '127.0.0.1', resolve);
});
servers.add(server);
return (server.address() as AddressInfo).port;
};
const closeServer = async (server: http.Server) => {
server.closeAllConnections();
await new Promise<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
servers.delete(server);
};
const withTimeout = async <T>(promise: Promise<T>, message: string): Promise<T> => {
let timer: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
promise,
new Promise<never>((_, reject) => {
timer = setTimeout(() => reject(new Error(message)), 1000);
})
]);
} finally {
if (timer) clearTimeout(timer);
}
};
afterEach(async () => {
await Promise.all([...servers].map(closeServer));
});
describe('MinIO timeout transport integration', () => {
it('destroys a real socket when the server never sends response headers', async () => {
const sockets = new Set<Socket>();
const server = http.createServer(() => {});
server.on('connection', (socket) => {
sockets.add(socket);
socket.once('close', () => sockets.delete(socket));
});
const port = await listen(server);
const transport = createMinioTimeoutTransport({ transport: http, timeoutMs: 50 });
const error = await new Promise<Error>((resolve, reject) => {
const request = transport.request({ host: '127.0.0.1', port, path: '/' });
request.once('error', resolve);
request.once('response', () => reject(new Error('Unexpected response')));
request.end();
});
expect(error.message).toBe('MinIO request timeout after 50ms');
await expect.poll(() => sockets.size, { timeout: 1000 }).toBe(0);
await closeServer(server);
});
it('destroys a real socket when the response body never completes', async () => {
const sockets = new Set<Socket>();
const server = http.createServer((_request, response) => {
response.writeHead(200, { 'Content-Length': '100' });
response.write('partial');
});
server.on('connection', (socket) => {
sockets.add(socket);
socket.once('close', () => sockets.delete(socket));
});
const port = await listen(server);
const transport = createMinioTimeoutTransport({ transport: http, timeoutMs: 50 });
let responseComplete: boolean | undefined;
await new Promise<void>((resolve, reject) => {
const request = transport.request({ host: '127.0.0.1', port, path: '/' }, (response) => {
response.resume();
response.once('end', () => reject(new Error('Unexpected complete response')));
response.once('error', () => {
responseComplete = response.complete;
resolve();
});
response.once('aborted', () => {
responseComplete = response.complete;
resolve();
});
});
request.once('error', () => {});
request.end();
});
expect(responseComplete).toBe(false);
await expect.poll(() => sockets.size, { timeout: 1000 }).toBe(0);
await closeServer(server);
});
it('closes a real download socket when the caller aborts after response headers', async () => {
const sockets = new Set<Socket>();
const server = http.createServer((_request, response) => {
response.writeHead(200, {
'Content-Length': '100',
'Content-Type': 'application/octet-stream'
});
response.write('partial');
});
server.on('connection', (socket) => {
sockets.add(socket);
socket.once('close', () => sockets.delete(socket));
});
const port = await listen(server);
const storage = new MinioStorageAdapter({
vendor: 'minio',
bucket: 'test-bucket',
endpoint: `http://127.0.0.1:${port}`,
region: 'us-east-1',
forcePathStyle: true,
maxRetries: 1,
credentials: { accessKeyId: 'access-key', secretAccessKey: 'secret-key' }
});
const controller = new AbortController();
const { body } = await storage.downloadObject({
key: 'abort/file.bin',
abortSignal: controller.signal
});
body.on('error', () => {});
const streamClosed = new Promise<void>((resolve) => body.once('close', resolve));
controller.abort(new Error('client aborted'));
await streamClosed;
expect(body.destroyed).toBe(true);
await expect.poll(() => sockets.size, { timeout: 1000 }).toBe(0);
await storage.destroy();
await closeServer(server);
});
it('aborts a real AWS-compatible request while waiting for response headers', async () => {
const sockets = new Set<Socket>();
let notifyRequest: (() => void) | undefined;
const requestReceived = new Promise<void>((resolve) => {
notifyRequest = resolve;
});
const server = http.createServer(() => notifyRequest?.());
server.on('connection', (socket) => {
sockets.add(socket);
socket.once('close', () => sockets.delete(socket));
});
const port = await listen(server);
const storage = new MinioStorageAdapter({
vendor: 'minio',
bucket: 'test-bucket',
endpoint: `http://127.0.0.1:${port}`,
region: 'us-east-1',
forcePathStyle: true,
maxRetries: 1,
credentials: { accessKeyId: 'access-key', secretAccessKey: 'secret-key' }
});
const controller = new AbortController();
const downloadPromise = storage.downloadObject({
key: 'abort/waiting-for-headers.bin',
abortSignal: controller.signal
});
await withTimeout(requestReceived, 'AWS-compatible request did not reach the local server');
controller.abort(new Error('client aborted'));
await expect(downloadPromise).rejects.toMatchObject({ name: 'AbortError' });
await expect.poll(() => sockets.size, { timeout: 1000 }).toBe(0);
await storage.destroy();
await closeServer(server);
});
});
import * as http from 'node:http';
import type { AddressInfo, Socket } from 'node:net';
import { afterEach, describe, expect, it } from 'vitest';
import { CosStorageAdapter } from '../../../src/adapters/cos.adapter';
import { OssStorageAdapter } from '../../../src/adapters/oss.adapter';
const servers = new Set<http.Server>();
const listen = async (server: http.Server): Promise<number> => {
await new Promise<void>((resolve, reject) => {
server.once('error', reject);
server.listen(0, '127.0.0.1', resolve);
});
servers.add(server);
return (server.address() as AddressInfo).port;
};
const closeServer = async (server: http.Server) => {
server.closeAllConnections();
await new Promise<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
servers.delete(server);
};
const createStalledObjectServer = () => {
const sockets = new Set<Socket>();
let notifyRequest: (() => void) | undefined;
const requestReceived = new Promise<void>((resolve) => {
notifyRequest = resolve;
});
const server = http.createServer((_request, response) => {
response.writeHead(200, {
'Content-Length': '100',
'Content-Type': 'application/octet-stream'
});
response.write('partial');
notifyRequest?.();
});
server.on('connection', (socket) => {
sockets.add(socket);
socket.once('close', () => sockets.delete(socket));
});
return { server, sockets, requestReceived };
};
const withTimeout = async <T>(promise: Promise<T>, message: string): Promise<T> => {
let timer: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
promise,
new Promise<never>((_, reject) => {
timer = setTimeout(() => reject(new Error(message)), 1000);
})
]);
} finally {
if (timer) clearTimeout(timer);
}
};
afterEach(async () => {
await Promise.all([...servers].map(closeServer));
});
describe('provider download abort transport integration', () => {
it('closes the OSS response socket when the caller aborts an in-flight stream', async () => {
const { server, sockets, requestReceived } = createStalledObjectServer();
const port = await listen(server);
const storage = new OssStorageAdapter({
vendor: 'oss',
bucket: 'test-bucket',
region: 'oss-cn-hangzhou',
endpoint: `http://127.0.0.1:${port}`,
cname: true,
secure: false,
credentials: { accessKeyId: 'access-key', secretAccessKey: 'secret-key' }
});
const controller = new AbortController();
const abortReason = new Error('client aborted');
const { body } = await storage.downloadObject({
key: 'abort/file.bin',
abortSignal: controller.signal
});
body.on('error', () => {});
const streamClosed = new Promise<void>((resolve) => body.once('close', resolve));
await withTimeout(requestReceived, 'OSS request did not reach the local server');
controller.abort(abortReason);
await withTimeout(streamClosed, 'OSS stream did not close after abort');
expect(body.errored).toBe(abortReason);
await expect.poll(() => sockets.size, { timeout: 1000 }).toBe(0);
await storage.destroy();
await closeServer(server);
});
it('aborts the COS request when the caller aborts its output stream', async () => {
const { server, sockets, requestReceived } = createStalledObjectServer();
const port = await listen(server);
const storage = new CosStorageAdapter({
vendor: 'cos',
bucket: 'test-bucket-1250000000',
region: 'ap-guangzhou',
protocol: 'http:',
domain: `127.0.0.1:${port}`,
credentials: { accessKeyId: 'secret-id', secretAccessKey: 'secret-key' }
});
const controller = new AbortController();
const abortReason = new Error('client aborted');
const { body } = await storage.downloadObject({
key: 'abort/file.bin',
abortSignal: controller.signal
});
body.on('error', () => {});
const streamClosed = new Promise<void>((resolve) => body.once('close', resolve));
await withTimeout(requestReceived, 'COS request did not reach the local server');
controller.abort(abortReason);
await withTimeout(streamClosed, 'COS stream did not close after abort');
expect(body.errored).toBe(abortReason);
await expect.poll(() => sockets.size, { timeout: 1000 }).toBe(0);
await storage.destroy();
await closeServer(server);
});
});
......@@ -8,7 +8,7 @@ import {
createS3AccessLinkService,
encodeExpiresAtMinute,
type S3DownloadAliasStore
} from '@fastgpt-sdk/storage/access-link';
} from '../../src/access-link';
const baseNow = new Date('2026-01-01T00:00:00.000Z');
const getFutureDate = (minutes: number) => new Date(baseNow.getTime() + minutes * 60_000);
......
import { PassThrough, Readable } from 'node:stream';
import { describe, expect, it, vi } from 'vitest';
import { AwsS3StorageAdapter } from '../../../src/adapters/aws-s3.adapter';
const createAdapter = () =>
new AwsS3StorageAdapter({
vendor: 'aws-s3',
bucket: 'fastgpt-private',
endpoint: 'http://localhost:9000',
region: 'us-east-1',
forcePathStyle: true,
maxRetries: 1,
credentials: {
accessKeyId: 'access-key',
secretAccessKey: 'secret-key'
}
});
describe('AwsS3StorageAdapter.downloadObject', () => {
it('rejects a pre-aborted download without dispatching an AWS request', async () => {
const adapter = createAdapter();
const send = vi.fn();
(adapter as any).client.send = send;
const controller = new AbortController();
controller.abort();
await expect(
adapter.downloadObject({
key: 'dataset/team/file.txt',
abortSignal: controller.signal
})
).rejects.toMatchObject({ name: 'AbortError' });
expect(send).not.toHaveBeenCalled();
});
it('passes the caller abort signal to the AWS request handler', async () => {
const adapter = createAdapter();
const body = Readable.from([Buffer.from('file')]);
const send = vi.fn().mockResolvedValue({ Body: body });
(adapter as any).client.send = send;
const controller = new AbortController();
const result = await adapter.downloadObject({
key: 'dataset/team/file.txt',
abortSignal: controller.signal
});
expect(result.body).toBe(body);
expect(send).toHaveBeenCalledWith(
expect.objectContaining({
input: {
Bucket: 'fastgpt-private',
Key: 'dataset/team/file.txt'
}
}),
{ abortSignal: controller.signal }
);
});
it('destroys an in-flight body with the caller abort reason', async () => {
const adapter = createAdapter();
const body = new PassThrough();
(adapter as any).client.send = vi.fn().mockResolvedValue({ Body: body });
const controller = new AbortController();
const abortReason = new Error('client aborted');
const result = await adapter.downloadObject({
key: 'dataset/team/file.txt',
abortSignal: controller.signal
});
result.body.on('error', () => {});
controller.abort(abortReason);
expect(result.body.errored).toBe(abortReason);
expect(result.body.destroyed).toBe(true);
});
});
describe('AwsS3StorageAdapter.deleteObjectsByPrefix', () => {
it('rejects a whitespace-only prefix without calling S3', async () => {
const adapter = createAdapter();
const send = vi.fn();
(adapter as any).client.send = send;
await expect(adapter.deleteObjectsByPrefix({ prefix: ' ' })).rejects.toThrow(
'Prefix is required'
);
expect(send).not.toHaveBeenCalled();
});
it('preserves failures collected before a later listing page is empty', async () => {
const adapter = createAdapter();
const send = vi
.fn()
.mockResolvedValueOnce({
Contents: [{ Key: 'dataset/failed.txt' }],
IsTruncated: true,
NextContinuationToken: 'next-page'
})
.mockResolvedValueOnce({ Errors: [{ Key: 'dataset/failed.txt' }] })
.mockResolvedValueOnce({ Contents: [], IsTruncated: false });
(adapter as any).client.send = send;
await expect(adapter.deleteObjectsByPrefix({ prefix: 'dataset/' })).resolves.toEqual({
bucket: 'fastgpt-private',
keys: ['dataset/failed.txt']
});
});
});
describe('AwsS3StorageAdapter.generatePublicGetUrl', () => {
it.each([
[
{ forcePathStyle: true, publicAccessExtraSubPath: undefined },
'https://storage.example.com/fastgpt-private/folder%20name/file%20%23%2B.txt'
],
[
{ forcePathStyle: true, publicAccessExtraSubPath: '/proxy/' },
'https://storage.example.com/proxy/fastgpt-private/folder%20name/file%20%23%2B.txt'
],
[
{ forcePathStyle: false, publicAccessExtraSubPath: undefined },
'https://fastgpt-private.storage.example.com/folder%20name/file%20%23%2B.txt'
],
[
{ forcePathStyle: false, publicAccessExtraSubPath: '/proxy/' },
'https://fastgpt-private.storage.example.com/proxy/folder%20name/file%20%23%2B.txt'
]
])('encodes keys for options %j', (overrides, expectedUrl) => {
const adapter = new AwsS3StorageAdapter({
vendor: 'aws-s3',
bucket: 'fastgpt-private',
endpoint: 'https://storage.example.com',
region: 'us-east-1',
credentials: {
accessKeyId: 'access-key',
secretAccessKey: 'secret-key'
},
...overrides
});
expect(adapter.generatePublicGetUrl({ key: 'folder name/file #+.txt' }).url).toBe(expectedUrl);
});
});
import { beforeEach, describe, expect, it, vi } from 'vitest';
import { CosStorageAdapter } from '../../../../sdk/storage/src/adapters/cos.adapter';
import { CosStorageAdapter } from '../../../src/adapters/cos.adapter';
const createAdapter = () =>
new CosStorageAdapter({
......@@ -52,6 +52,22 @@ describe('CosStorageAdapter.generatePresignedGetUrl', () => {
});
describe('CosStorageAdapter.downloadObject', () => {
it('rejects a pre-aborted download without requesting the object', async () => {
const adapter = createAdapter();
const getObject = vi.fn();
(adapter as any).client.getObject = getObject;
const controller = new AbortController();
controller.abort();
await expect(
adapter.downloadObject({
key: 'dataset/team/file.txt',
abortSignal: controller.signal
})
).rejects.toMatchObject({ name: 'AbortError' });
expect(getObject).not.toHaveBeenCalled();
});
it('destroys the output stream when the caller aborts the download', async () => {
const adapter = createAdapter();
(adapter as any).client.getObject = vi.fn();
......@@ -61,8 +77,57 @@ describe('CosStorageAdapter.downloadObject', () => {
key: 'dataset/team/file.txt',
abortSignal: controller.signal
});
controller.abort(new Error('client aborted'));
const abortReason = new Error('client aborted');
body.on('error', () => {});
controller.abort(abortReason);
expect(body.errored).toBe(abortReason);
expect(body.destroyed).toBe(true);
});
});
describe('CosStorageAdapter deletion boundaries', () => {
it('treats an empty key list as a no-op', async () => {
const adapter = createAdapter();
const deleteMultipleObject = vi.fn();
(adapter as any).client.deleteMultipleObject = deleteMultipleObject;
await expect(adapter.deleteObjectsByMultiKeys({ keys: [] })).resolves.toEqual({
bucket: 'fastgpt-private',
keys: []
});
expect(deleteMultipleObject).not.toHaveBeenCalled();
});
it('rejects a whitespace-only prefix without listing objects', async () => {
const adapter = createAdapter();
const getBucket = vi.fn();
(adapter as any).client.getBucket = getBucket;
await expect(adapter.deleteObjectsByPrefix({ prefix: ' ' })).rejects.toThrow(
'Prefix is required'
);
expect(getBucket).not.toHaveBeenCalled();
});
});
describe('CosStorageAdapter.generatePublicGetUrl', () => {
it.each([
[undefined, 'https://fastgpt-private.cos.ap-guangzhou.myqcloud.com/folder%20%23/file%2B.txt'],
['cdn.example.com', 'https://cdn.example.com/folder%20%23/file%2B.txt']
])('encodes keys with domain %j', (domain, expectedUrl) => {
const adapter = new CosStorageAdapter({
vendor: 'cos',
bucket: 'fastgpt-private',
region: 'ap-guangzhou',
protocol: 'https:',
domain,
credentials: {
accessKeyId: 'secret-id',
secretAccessKey: 'secret-key'
}
});
expect(adapter.generatePublicGetUrl({ key: 'folder #/file+.txt' }).url).toBe(expectedUrl);
});
});
import { PassThrough } from 'node:stream';
import { describe, expect, it, vi } from 'vitest';
import { OssStorageAdapter } from '../../../src/adapters/oss.adapter';
const createAdapter = () =>
new OssStorageAdapter({
vendor: 'oss',
bucket: 'fastgpt-private',
endpoint: 'http://localhost:9000',
region: 'oss-cn-hangzhou',
secure: false,
credentials: {
accessKeyId: 'access-key',
secretAccessKey: 'secret-key'
}
});
describe('OssStorageAdapter.downloadObject', () => {
it('rejects a pre-aborted download without requesting a stream', async () => {
const adapter = createAdapter();
const getStream = vi.fn();
(adapter as any).client.getStream = getStream;
const controller = new AbortController();
controller.abort();
await expect(
adapter.downloadObject({ key: 'dataset/file.txt', abortSignal: controller.signal })
).rejects.toMatchObject({ name: 'AbortError' });
expect(getStream).not.toHaveBeenCalled();
});
it('destroys an in-flight stream with the caller abort reason', async () => {
const adapter = createAdapter();
const stream = new PassThrough();
(adapter as any).client.getStream = vi.fn().mockResolvedValue({ stream });
const controller = new AbortController();
const abortReason = new Error('client aborted');
const result = await adapter.downloadObject({
key: 'dataset/file.txt',
abortSignal: controller.signal
});
result.body.on('error', () => {});
controller.abort(abortReason);
expect(result.body.errored).toBe(abortReason);
expect(result.body.destroyed).toBe(true);
});
});
describe('OssStorageAdapter deletion boundaries', () => {
it('treats an empty key list as a no-op', async () => {
const adapter = createAdapter();
const deleteMulti = vi.fn();
(adapter as any).client.deleteMulti = deleteMulti;
await expect(adapter.deleteObjectsByMultiKeys({ keys: [] })).resolves.toEqual({
bucket: 'fastgpt-private',
keys: []
});
expect(deleteMulti).not.toHaveBeenCalled();
});
it('returns only keys missing from the verbose delete response as failures', async () => {
const adapter = createAdapter();
const deleteMulti = vi.fn().mockResolvedValue({ deleted: ['first.txt'] });
(adapter as any).client.deleteMulti = deleteMulti;
await expect(
adapter.deleteObjectsByMultiKeys({ keys: ['first.txt', 'second.txt'] })
).resolves.toEqual({
bucket: 'fastgpt-private',
keys: ['second.txt']
});
expect(deleteMulti).toHaveBeenCalledWith(['first.txt', 'second.txt'], { quiet: false });
});
it('normalizes the object-shaped Deleted entries produced by ali-oss XML parsing', async () => {
const adapter = createAdapter();
(adapter as any).client.deleteMulti = vi.fn().mockResolvedValue({
deleted: [{ Key: 'first.txt' }, { Key: 'second.txt', VersionId: 'version-1' }]
});
await expect(
adapter.deleteObjectsByMultiKeys({ keys: ['first.txt', 'second.txt'] })
).resolves.toEqual({
bucket: 'fastgpt-private',
keys: []
});
});
it('conservatively returns every key when the verbose response omits deleted entries', async () => {
const adapter = createAdapter();
(adapter as any).client.deleteMulti = vi.fn().mockResolvedValue({});
await expect(
adapter.deleteObjectsByMultiKeys({ keys: ['first.txt', 'second.txt'] })
).resolves.toEqual({
bucket: 'fastgpt-private',
keys: ['first.txt', 'second.txt']
});
});
it('preserves failures collected before a later listing page is empty', async () => {
const adapter = createAdapter();
(adapter as any).client.list = vi
.fn()
.mockResolvedValueOnce({
objects: [{ name: 'dataset/failed.txt' }],
isTruncated: true,
nextMarker: 'dataset/next.txt'
})
.mockResolvedValueOnce({ objects: [], isTruncated: false });
(adapter as any).client.deleteMulti = vi.fn().mockResolvedValue({ deleted: [] });
await expect(adapter.deleteObjectsByPrefix({ prefix: 'dataset/' })).resolves.toEqual({
bucket: 'fastgpt-private',
keys: ['dataset/failed.txt']
});
});
it('rejects a whitespace-only prefix without listing objects', async () => {
const adapter = createAdapter();
const list = vi.fn();
(adapter as any).client.list = list;
await expect(adapter.deleteObjectsByPrefix({ prefix: ' ' })).rejects.toThrow(
'Prefix is required'
);
expect(list).not.toHaveBeenCalled();
});
});
describe('OssStorageAdapter.generatePublicGetUrl', () => {
it.each([
[
false,
undefined,
'https://fastgpt-private.oss-cn-hangzhou.aliyuncs.com/folder%20%23/file%2B.txt'
],
[true, 'cdn.example.com', 'https://cdn.example.com/folder%20%23/file%2B.txt']
])('encodes keys with cname=%s', (cname, endpoint, expectedUrl) => {
const adapter = new OssStorageAdapter({
vendor: 'oss',
bucket: 'fastgpt-private',
endpoint,
region: 'oss-cn-hangzhou',
secure: true,
cname,
credentials: {
accessKeyId: 'access-key',
secretAccessKey: 'secret-key'
}
});
expect(adapter.generatePublicGetUrl({ key: 'folder #/file+.txt' }).url).toBe(expectedUrl);
});
});
import { describe, expect, it, vi } from 'vitest';
import type { IStorage } from '../../src/interface';
import { AwsS3StorageAdapter } from '../../src/adapters/aws-s3.adapter';
import { CosStorageAdapter } from '../../src/adapters/cos.adapter';
import { MinioStorageAdapter } from '../../src/adapters/minio.adapter';
import { OssStorageAdapter } from '../../src/adapters/oss.adapter';
import { InvalidStorageObjectKeyError } from '../../src/errors';
import { createVitestStorageMock } from '../../src/helper/mock';
type GuardedStorage = {
storage: IStorage;
remoteCall: ReturnType<typeof vi.fn>;
};
const unexpectedRemoteCall = () => {
throw new Error('Object key was not validated before the remote SDK call');
};
const createAwsStorage = (): GuardedStorage => {
const storage = new AwsS3StorageAdapter({
vendor: 'aws-s3',
bucket: 'test-bucket',
endpoint: 'http://127.0.0.1:1',
region: 'us-east-1',
forcePathStyle: true,
maxRetries: 1,
credentials: { accessKeyId: 'access-key', secretAccessKey: 'secret-key' }
});
const remoteCall = vi.fn(unexpectedRemoteCall);
(storage as any).client.send = remoteCall;
return { storage, remoteCall };
};
const createMinioStorage = (): GuardedStorage => {
const storage = new MinioStorageAdapter({
vendor: 'minio',
bucket: 'test-bucket',
endpoint: 'http://127.0.0.1:1',
region: 'us-east-1',
forcePathStyle: true,
maxRetries: 1,
credentials: { accessKeyId: 'access-key', secretAccessKey: 'secret-key' }
});
const remoteCall = vi.fn(unexpectedRemoteCall);
(storage as any).client.send = remoteCall;
(storage as any).minioClient.removeObject = remoteCall;
(storage as any).minioClient.removeObjects = remoteCall;
return { storage, remoteCall };
};
const createOssStorage = (): GuardedStorage => {
const storage = new OssStorageAdapter({
vendor: 'oss',
bucket: 'test-bucket',
endpoint: 'http://127.0.0.1:1',
region: 'oss-cn-hangzhou',
secure: false,
credentials: { accessKeyId: 'access-key', secretAccessKey: 'secret-key' }
});
const remoteCall = vi.fn(unexpectedRemoteCall);
Object.assign((storage as any).client, {
head: remoteCall,
put: remoteCall,
getStream: remoteCall,
delete: remoteCall,
deleteMulti: remoteCall,
list: remoteCall,
signatureUrlV4: remoteCall,
signatureUrl: remoteCall,
copy: remoteCall
});
return { storage, remoteCall };
};
const createCosStorage = (): GuardedStorage => {
const storage = new CosStorageAdapter({
vendor: 'cos',
bucket: 'test-bucket',
region: 'ap-guangzhou',
credentials: { accessKeyId: 'access-key', secretAccessKey: 'secret-key' }
});
const remoteCall = vi.fn(unexpectedRemoteCall);
Object.assign((storage as any).client, {
headObject: remoteCall,
putObject: remoteCall,
getObject: remoteCall,
deleteObject: remoteCall,
deleteMultipleObject: remoteCall,
getBucket: remoteCall,
getObjectUrl: remoteCall,
sliceCopyFile: remoteCall
});
return { storage, remoteCall };
};
const createMockStorage = (): GuardedStorage => ({
storage: createVitestStorageMock({ vi }),
remoteCall: vi.fn()
});
const storageFactories = [
['AWS S3', createAwsStorage],
['MinIO', createMinioStorage],
['OSS', createOssStorage],
['COS', createCosStorage],
['Vitest mock', createMockStorage]
] as const;
const invalidKey = 'invalid//key';
const keyOperations: ReadonlyArray<
[name: string, field: string, operation: (storage: IStorage) => unknown]
> = [
['checkObjectExists', 'key', (storage) => storage.checkObjectExists({ key: invalidKey })],
['getObjectMetadata', 'key', (storage) => storage.getObjectMetadata({ key: invalidKey })],
['uploadObject', 'key', (storage) => storage.uploadObject({ key: invalidKey, body: 'body' })],
['downloadObject', 'key', (storage) => storage.downloadObject({ key: invalidKey })],
['deleteObject', 'key', (storage) => storage.deleteObject({ key: invalidKey })],
[
'deleteObjectsByMultiKeys',
'keys[1]',
(storage) => storage.deleteObjectsByMultiKeys({ keys: ['valid/key', invalidKey] })
],
[
'deleteObjectsByPrefix',
'prefix',
(storage) => storage.deleteObjectsByPrefix({ prefix: invalidKey })
],
[
'generatePresignedPutUrl',
'key',
(storage) => storage.generatePresignedPutUrl({ key: invalidKey })
],
[
'generatePresignedGetUrl',
'key',
(storage) => storage.generatePresignedGetUrl({ key: invalidKey })
],
['generatePublicGetUrl', 'key', (storage) => storage.generatePublicGetUrl({ key: invalidKey })],
['listObjects', 'prefix', (storage) => storage.listObjects({ prefix: invalidKey })],
[
'copyObjectInSelfBucket source',
'sourceKey',
(storage) =>
storage.copyObjectInSelfBucket({ sourceKey: invalidKey, targetKey: 'valid/target' })
],
[
'copyObjectInSelfBucket target',
'targetKey',
(storage) =>
storage.copyObjectInSelfBucket({ sourceKey: 'valid/source', targetKey: invalidKey })
]
];
type KeyContractCase = readonly [
storageName: string,
operationName: string,
field: string,
createStorage: () => GuardedStorage,
operation: (storage: IStorage) => unknown
];
const keyContractCases: KeyContractCase[] = storageFactories.flatMap(
([storageName, createStorage]) =>
keyOperations.map<KeyContractCase>(([operationName, field, operation]) => [
storageName,
operationName,
field,
createStorage,
operation
])
);
describe('IStorage object key preflight contract', () => {
it.each(keyContractCases)(
'%s validates %s before dispatch',
async (_storageName, _operationName, field, createStorage, operation) => {
const { storage, remoteCall } = createStorage();
await expect(Promise.resolve().then(() => operation(storage))).rejects.toMatchObject({
name: InvalidStorageObjectKeyError.name,
field,
reason: 'empty_path_segment'
});
expect(remoteCall).not.toHaveBeenCalled();
}
);
it.each([
['AWS S3', createAwsStorage],
['MinIO', createMinioStorage]
] as const)(
'%s validates every key before deleting the first 1000-key batch',
async (_, createStorage) => {
const { storage, remoteCall } = createStorage();
const keys = Array.from({ length: 1000 }, (_, index) => `valid/${index}`).concat(invalidKey);
await expect(storage.deleteObjectsByMultiKeys({ keys })).rejects.toBeInstanceOf(
InvalidStorageObjectKeyError
);
expect(remoteCall).not.toHaveBeenCalled();
}
);
});
import { describe, expect, it } from 'vitest';
import { InvalidStorageObjectKeyError } from '../../src/errors';
import {
MAX_STORAGE_OBJECT_KEY_UTF8_BYTES,
assertRequiredStorageObjectPrefix,
assertStorageObjectKey,
assertStorageObjectKeys,
assertStorageObjectPrefix
} from '../../src/assert';
const expectInvalidKey = ({
value,
reason,
field = 'key'
}: {
value: unknown;
reason: InvalidStorageObjectKeyError['reason'];
field?: string;
}) => {
try {
assertStorageObjectKey(value, field);
throw new Error('Expected key validation to fail');
} catch (error) {
expect(error).toBeInstanceOf(InvalidStorageObjectKeyError);
expect(error).toMatchObject({ reason, field });
}
};
describe('storage object key validation', () => {
it.each([
'folder name/+ # & % ?/\u6587\u4ef6-\ud83d\ude00.txt',
'folder/.hidden',
'folder/..backup',
'folder/trailing/'
])('accepts portable URL-sensitive and Unicode key %j', (key) => {
expect(() => assertStorageObjectKey(key)).not.toThrow();
});
it('accepts exactly 850 UTF-8 bytes, including multibyte characters', () => {
const asciiKey = 'a'.repeat(MAX_STORAGE_OBJECT_KEY_UTF8_BYTES);
const unicodeKey = `${'\u4e2d'.repeat(282)}abcd`;
expect(Buffer.byteLength(asciiKey)).toBe(850);
expect(Buffer.byteLength(unicodeKey)).toBe(850);
expect(() => assertStorageObjectKey(asciiKey)).not.toThrow();
expect(() => assertStorageObjectKey(unicodeKey)).not.toThrow();
});
it('rejects 851 UTF-8 bytes even when the JavaScript string is shorter', () => {
const key = `${'\u4e2d'.repeat(283)}ab`;
expect(key.length).toBeLessThan(850);
expect(Buffer.byteLength(key)).toBe(851);
expectInvalidKey({ value: key, reason: 'too_long' });
});
it.each([
[undefined, 'invalid_type'],
[123, 'invalid_type'],
['', 'empty'],
['bad\ud800key', 'invalid_unicode'],
['bad\udcffkey', 'invalid_unicode'],
['/leading', 'leading_slash'],
['folder\\file', 'backslash'],
['folder//file', 'empty_path_segment'],
['folder\0file', 'control_character'],
['folder\nfile', 'control_character'],
['folder\u007ffile', 'control_character']
] as const)('rejects non-portable key %j as %s', (value, reason) => {
expectInvalidKey({ value, reason });
});
it.each(['.', '..', './file', 'folder/./file', 'folder/../file', 'folder/ .. /file'])(
'rejects dot path segment in %j',
(key) => {
expectInvalidKey({ value: key, reason: 'dot_path_segment' });
}
);
it.each(['\u0018', '\u0019', '\u001a', '\u001b'])(
'rejects COS-incompatible control character U+%s',
(character) => {
expectInvalidKey({ value: `folder/${character}/file`, reason: 'control_character' });
}
);
it('reports the first invalid batch item and validates the full batch before callers mutate', () => {
try {
assertStorageObjectKeys(['valid/first', 'invalid//second', 'valid/third']);
throw new Error('Expected key validation to fail');
} catch (error) {
expect(error).toBeInstanceOf(InvalidStorageObjectKeyError);
expect(error).toMatchObject({
field: 'keys[1]',
reason: 'empty_path_segment'
});
}
});
it('rejects a non-array batch with the structured key error', () => {
try {
assertStorageObjectKeys('not-an-array');
throw new Error('Expected key validation to fail');
} catch (error) {
expect(error).toBeInstanceOf(InvalidStorageObjectKeyError);
expect(error).toMatchObject({ field: 'keys', reason: 'invalid_type' });
}
});
it('rejects sparse arrays instead of skipping missing entries', () => {
const keys = new Array<string>(2);
keys[1] = 'valid.txt';
try {
assertStorageObjectKeys(keys);
throw new Error('Expected key validation to fail');
} catch (error) {
expect(error).toBeInstanceOf(InvalidStorageObjectKeyError);
expect(error).toMatchObject({ field: 'keys[0]', reason: 'invalid_type' });
}
});
it('allows omitted and empty list prefixes but validates non-empty prefixes', () => {
expect(() => assertStorageObjectPrefix(undefined)).not.toThrow();
expect(() => assertStorageObjectPrefix('')).not.toThrow();
expect(() => assertStorageObjectPrefix('folder name/\u4e2d\u6587/')).not.toThrow();
try {
assertStorageObjectPrefix('invalid//prefix');
throw new Error('Expected prefix validation to fail');
} catch (error) {
expect(error).toBeInstanceOf(InvalidStorageObjectKeyError);
expect(error).toMatchObject({ field: 'prefix', reason: 'empty_path_segment' });
}
});
it.each(['', ' '])('rejects required delete prefix %j', (prefix) => {
expect(() => assertRequiredStorageObjectPrefix(prefix)).toThrow('Prefix is required');
});
it('accepts a valid required delete prefix', () => {
expect(() => assertRequiredStorageObjectPrefix('team/files/')).not.toThrow();
});
});
import { describe, expect, it } from 'vitest';
import { AwsS3StorageAdapter } from '../../src/adapters/aws-s3.adapter';
import { CosStorageAdapter } from '../../src/adapters/cos.adapter';
import { MinioStorageAdapter } from '../../src/adapters/minio.adapter';
import { OssStorageAdapter } from '../../src/adapters/oss.adapter';
import { createStorage } from '../../src/factory';
const credentials = {
accessKeyId: 'access-key',
secretAccessKey: 'secret-key'
};
describe('createStorage', () => {
it('creates every supported adapter from its discriminated options', () => {
expect(
createStorage({
vendor: 'aws-s3',
bucket: 'bucket',
endpoint: 'http://localhost:9000',
region: 'us-east-1',
credentials
})
).toBeInstanceOf(AwsS3StorageAdapter);
expect(
createStorage({
vendor: 'minio',
bucket: 'bucket',
endpoint: 'http://localhost:9000',
region: 'us-east-1',
credentials
})
).toBeInstanceOf(MinioStorageAdapter);
expect(
createStorage({
vendor: 'oss',
bucket: 'bucket',
region: 'oss-cn-hangzhou',
credentials
})
).toBeInstanceOf(OssStorageAdapter);
expect(
createStorage({
vendor: 'cos',
bucket: 'bucket',
region: 'ap-guangzhou',
credentials
})
).toBeInstanceOf(CosStorageAdapter);
});
it('rejects an unsupported runtime vendor', () => {
expect(() => createStorage({ vendor: 'unknown' } as never)).toThrow(
'Unsupported storage vendor: unknown'
);
});
});
import { describe, expect, it, vi } from 'vitest';
import { createVitestStorageMock } from '../../src/helper/mock';
import { removeIntegrationBucketIfExists } from '../integration/helpers';
describe('removeIntegrationBucketIfExists', () => {
it('does nothing when the stable test bucket does not exist', async () => {
const storage = createVitestStorageMock({ vi });
const bucketExists = vi.fn().mockResolvedValue(false);
const deleteBucket = vi.fn();
await removeIntegrationBucketIfExists({ storage, bucketExists, deleteBucket });
expect(storage.listObjects).not.toHaveBeenCalled();
expect(deleteBucket).not.toHaveBeenCalled();
});
it('deletes an existing empty bucket', async () => {
const storage = createVitestStorageMock({ vi });
const deleteBucket = vi.fn().mockResolvedValue(undefined);
await removeIntegrationBucketIfExists({
storage,
bucketExists: vi.fn().mockResolvedValue(true),
deleteBucket
});
expect(storage.listObjects).toHaveBeenCalledWith({});
expect(storage.deleteObjectsByMultiKeys).not.toHaveBeenCalled();
expect(deleteBucket).toHaveBeenCalledOnce();
});
it('clears every object before deleting an existing bucket', async () => {
const storage = createVitestStorageMock({ vi });
storage.__putObject('first.txt', { body: Buffer.from('first') });
storage.__putObject('nested/second.txt', { body: Buffer.from('second') });
const deleteBucket = vi.fn().mockResolvedValue(undefined);
await removeIntegrationBucketIfExists({
storage,
bucketExists: vi.fn().mockResolvedValue(true),
deleteBucket
});
expect(storage.deleteObjectsByMultiKeys).toHaveBeenCalledWith({
keys: ['first.txt', 'nested/second.txt']
});
expect(storage.__objects.size).toBe(0);
expect(deleteBucket).toHaveBeenCalledOnce();
});
it('keeps the bucket when any object deletion fails', async () => {
const storage = createVitestStorageMock({ vi });
storage.__putObject('failed.txt', { body: Buffer.from('failed') });
vi.spyOn(storage, 'deleteObjectsByMultiKeys').mockResolvedValue({
bucket: storage.bucketName,
keys: ['failed.txt']
});
const deleteBucket = vi.fn();
await expect(
removeIntegrationBucketIfExists({
storage,
bucketExists: vi.fn().mockResolvedValue(true),
deleteBucket
})
).rejects.toThrow('Failed to clean integration test bucket: failed.txt');
expect(deleteBucket).not.toHaveBeenCalled();
});
it('retries the same bucket on the next setup after bucket deletion fails', async () => {
const storage = createVitestStorageMock({ vi });
const bucketExists = vi.fn().mockResolvedValue(true);
const deleteBucket = vi
.fn()
.mockRejectedValueOnce(new Error('temporary delete failure'))
.mockResolvedValueOnce(undefined);
await expect(
removeIntegrationBucketIfExists({ storage, bucketExists, deleteBucket })
).rejects.toThrow('temporary delete failure');
await expect(
removeIntegrationBucketIfExists({ storage, bucketExists, deleteBucket })
).resolves.toBeUndefined();
expect(bucketExists).toHaveBeenCalledTimes(2);
expect(deleteBucket).toHaveBeenCalledTimes(2);
});
});
import { Readable } from 'node:stream';
import { describe, expect, it, vi } from 'vitest';
import { createVitestStorageMock } from '../../src/helper/mock';
const readBody = async (body: Readable) => {
const chunks: Buffer[] = [];
for await (const chunk of body) {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
}
return Buffer.concat(chunks);
};
describe('createVitestStorageMock', () => {
it('implements bucket lifecycle, upload forms, download and metadata', async () => {
const storage = createVitestStorageMock({ vi, bucketName: 'test-bucket' });
await expect(storage.ensureBucket()).resolves.toEqual({
bucket: 'test-bucket',
exists: false,
created: true
});
await expect(storage.ensureBucket()).resolves.toEqual({
bucket: 'test-bucket',
exists: true,
created: false
});
await storage.uploadObject({
key: 'buffer.bin',
body: Buffer.from([0, 255]),
metadata: { source: 'buffer' }
});
await storage.uploadObject({ key: 'string.txt', body: 'string' });
await storage.uploadObject({ key: 'stream.txt', body: Readable.from(['stream']) });
const download = await storage.downloadObject({ key: 'buffer.bin' });
await expect(readBody(download.body)).resolves.toEqual(Buffer.from([0, 255]));
await expect(storage.getObjectMetadata({ key: 'buffer.bin' })).resolves.toMatchObject({
metadata: { source: 'buffer' },
contentLength: 2,
etag: expect.any(String)
});
await expect(storage.listObjects({})).resolves.toMatchObject({
keys: ['buffer.bin', 'stream.txt', 'string.txt']
});
});
it('returns failed keys rather than deleted keys from successful deletion methods', async () => {
const storage = createVitestStorageMock({ vi });
storage.__putObject('prefix/first.txt', { body: Buffer.from('first') });
storage.__putObject('prefix/second.txt', { body: Buffer.from('second') });
await expect(
storage.deleteObjectsByMultiKeys({ keys: ['prefix/first.txt', 'missing.txt'] })
).resolves.toEqual({ bucket: 'mock-bucket', keys: [] });
await expect(storage.deleteObjectsByPrefix({ prefix: 'prefix/' })).resolves.toEqual({
bucket: 'mock-bucket',
keys: []
});
expect(storage.__objects.size).toBe(0);
});
it('rejects empty prefixes and pre-aborted downloads', async () => {
const storage = createVitestStorageMock({ vi });
storage.__putObject('file.txt', { body: Buffer.from('file') });
await expect(storage.deleteObjectsByPrefix({ prefix: ' ' })).rejects.toThrow(
'Prefix is required'
);
const controller = new AbortController();
controller.abort();
await expect(
storage.downloadObject({ key: 'file.txt', abortSignal: controller.signal })
).rejects.toMatchObject({ name: 'AbortError' });
});
it('destroys an in-flight download with the caller abort reason', async () => {
const storage = createVitestStorageMock({ vi });
storage.__putObject('file.txt', { body: Buffer.from('file') });
const controller = new AbortController();
const abortReason = new Error('client aborted');
const { body } = await storage.downloadObject({
key: 'file.txt',
abortSignal: controller.signal
});
body.on('error', () => {});
controller.abort(abortReason);
expect(body.errored).toBe(abortReason);
expect(body.destroyed).toBe(true);
});
it('copies object buffers independently and resets state', async () => {
const storage = createVitestStorageMock({ vi });
storage.__putObject('source.txt', {
body: Buffer.from('source'),
metadata: { copied: 'true' }
});
await storage.copyObjectInSelfBucket({
sourceKey: 'source.txt',
targetKey: 'target.txt'
});
storage.__objects.get('source.txt')?.body.fill(0);
expect(storage.__objects.get('target.txt')?.body.toString()).toBe('source');
expect(storage.__objects.get('target.txt')?.metadata).toEqual({ copied: 'true' });
storage.__reset();
expect(storage.__objects.size).toBe(0);
});
it('encodes keys and response overrides in generated URLs', async () => {
const storage = createVitestStorageMock({ vi, baseUrl: 'https://storage.test' });
const key = 'folder name/file#+.txt';
await expect(storage.generatePresignedPutUrl({ key })).resolves.toMatchObject({
url: `https://storage.test/put/mock-bucket/${encodeURIComponent(key)}`
});
await expect(
storage.generatePresignedGetUrl({ key, responseContentType: 'text/plain; charset=utf-8' })
).resolves.toMatchObject({
url: expect.stringContaining('response-content-type=text%2Fplain%3B%20charset%3Dutf-8')
});
expect(storage.generatePublicGetUrl({ key }).url).toBe(
`https://storage.test/public/mock-bucket/${encodeURIComponent(key)}`
);
});
});
import { describe, expect, it } from 'vitest';
import { encodeObjectKeyPath } from '../../src/utils';
describe('encodeObjectKeyPath', () => {
it.each([
['', ''],
['folder/file.txt', 'folder/file.txt'],
['folder name/file #+&.txt', 'folder%20name/file%20%23%2B%26.txt'],
['literal/%2F/\u6587\u4ef6.txt', 'literal/%252F/%E6%96%87%E4%BB%B6.txt'],
['/leading//empty/', '/leading//empty/']
])('encodes %j as an object URL path', (key, expected) => {
expect(encodeObjectKeyPath(key)).toBe(expected);
});
});
{
"extends": "./tsconfig.json",
"compilerOptions": {
"allowJs": false,
"noEmit": true
},
"include": ["src/**/*.ts", "test/**/*.ts", "vitest.config.ts"]
}
import path from 'node:path';
import { defineConfig } from 'vitest/config';
export default defineConfig(() => {
return {
test: {
globalSetup: path.join(import.meta.dirname, 'test/global-setup.ts'),
include: ['test/**/*.test.ts'],
pool: 'threads',
fileParallelism: false,
testTimeout: 120_000,
hookTimeout: 120_000,
reporters: ['default']
}
};
});
import { vi } from 'vitest';
import { createVitestStorageMock } from '../../../sdk/storage/src/testing/vitestMock';
import { createVitestStorageMock } from '../../../sdk/storage/src/helper/mock';
const mockStorageByBucket = new Map<string, ReturnType<typeof createVitestStorageMock>>();
const getMockStorage = (bucketName: string) => {
......
import { Readable } from 'node:stream';
import { describe, expect, it, vi } from 'vitest';
import { AwsS3StorageAdapter } from '../../../../sdk/storage/src/adapters/aws-s3.adapter';
const createAdapter = () =>
new AwsS3StorageAdapter({
vendor: 'aws-s3',
bucket: 'fastgpt-private',
endpoint: 'http://localhost:9000',
region: 'us-east-1',
forcePathStyle: true,
maxRetries: 1,
credentials: {
accessKeyId: 'access-key',
secretAccessKey: 'secret-key'
}
});
describe('AwsS3StorageAdapter.downloadObject', () => {
it('passes the caller abort signal to the AWS request handler', async () => {
const adapter = createAdapter();
const body = Readable.from([Buffer.from('file')]);
const send = vi.fn().mockResolvedValue({ Body: body });
(adapter as any).client.send = send;
const controller = new AbortController();
const result = await adapter.downloadObject({
key: 'dataset/team/file.txt',
abortSignal: controller.signal
});
expect(result.body).toBe(body);
expect(send).toHaveBeenCalledWith(
expect.objectContaining({
input: {
Bucket: 'fastgpt-private',
Key: 'dataset/team/file.txt'
}
}),
{ abortSignal: controller.signal }
);
});
});
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or sign in to comment