diff --git a/.changeset/config.json b/.changeset/config.json index 4f96e85b..00062adf 100644 --- a/.changeset/config.json +++ b/.changeset/config.json @@ -23,7 +23,9 @@ "@cleverbrush/server", "@cleverbrush/server-openapi", "@cleverbrush/orm", - "@cleverbrush/orm-cli" + "@cleverbrush/orm-cli", + "@cleverbrush/storage", + "@cleverbrush/storage-s3" ] ], "linked": [], diff --git a/.changeset/object-storage-s3.md b/.changeset/object-storage-s3.md new file mode 100644 index 00000000..a3aa19cc --- /dev/null +++ b/.changeset/object-storage-s3.md @@ -0,0 +1,13 @@ +--- +"@cleverbrush/storage": minor +"@cleverbrush/storage-s3": minor +--- + +Add provider-neutral object storage contracts, portable errors and safe public URL +mapping. Add an S3-compatible adapter with explicit custom endpoints, credentials, +addressing style, key prefixes and independently configured public asset URLs. +Support streamed reads, bounded multipart writes, metadata-preserving copies, +idempotent deletion, cancellation and asynchronous disposal through `await using` +on the shared storage contract. Include shared contract +coverage, real Garage integration tests, and configuration examples for hosted +and self-hosted services. diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2819da0b..27819052 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -76,3 +76,16 @@ jobs: - run: npm run test:scheduler:integration - run: node demos/durable-jobs/demo.ts - run: node demos/durable-jobs/periodic.ts --fast + + storage-integration: + name: S3 Storage Integration (Garage) + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-node@v4 + with: + node-version: 24 + cache: npm + - run: npm ci + - run: npm run build + - run: npm run test:storage:integration diff --git a/README.md b/README.md index 60bdae01..bc1a71c9 100644 --- a/README.md +++ b/README.md @@ -41,6 +41,8 @@ JSON Schema, API contracts, and Standard Schema integrations. | [`@cleverbrush/deep`](./libs/deep) | Deep equality, deep extension, flattening, and object utilities. | | [`@cleverbrush/scheduler`](./libs/scheduler) | Typed durable jobs, recurring triggers and ordered progress. | | [`@cleverbrush/scheduler-postgres`](./libs/scheduler-postgres) | PostgreSQL job persistence, transactional enqueue and fenced leases. | +| [`@cleverbrush/storage`](./libs/storage) | Provider-neutral object storage contracts and public URL mapping. | +| [`@cleverbrush/storage-s3`](./libs/storage-s3) | Streaming S3-compatible storage for self-hosted and hosted providers. | ## How The Pieces Fit diff --git a/docs/framework-feature-candidates.md b/docs/framework-feature-candidates.md index ff51e844..2039415f 100644 --- a/docs/framework-feature-candidates.md +++ b/docs/framework-feature-candidates.md @@ -1,6 +1,6 @@ # Framework feature candidates -Status: F01–F03 implemented on the feature branch for PR review; F04–F07 remain proposed. +Status: F01–F03 merged; F04–F05 implemented on the storage feature branch for PR review; F06–F07 remain proposed. Assessment date: 2026-10-02. This is an unprioritized list of reusable Framework capabilities and correctness @@ -163,6 +163,8 @@ application responsibilities. - Public URL construction handles object keys correctly and does not expose credentials or depend on temporary signed URLs. +**Implementation:** Added the provider-neutral storage contract, portable errors, key and public URL helpers, and a shared adapter contract suite. See the [storage guide](../libs/storage/README.md). + **Review notes:** ## F05 — S3-compatible storage adapter @@ -194,6 +196,8 @@ features are separate future candidates. - Verify object contents with an independent checksum or byte comparison rather than assuming an ETag always represents a content checksum. +**Implementation:** Added the configurable S3 adapter, bounded multipart transfers, cancellation and stream cleanup, Garage integration tests and a dedicated CI job. Hetzner configuration is documented; a live Hetzner account is not included in CI. See the [S3 guide](../libs/storage-s3/README.md). + **Review notes:** ## F06 — CORS preflight support diff --git a/libs/storage-s3/README.md b/libs/storage-s3/README.md new file mode 100644 index 00000000..9b9bc4f2 --- /dev/null +++ b/libs/storage-s3/README.md @@ -0,0 +1,168 @@ +# @cleverbrush/storage-s3 + +An implementation of [`ObjectStorage`](../storage) for S3-compatible services. +Uses the modular AWS SDK as a protocol client; an AWS account is not required. +Explicit endpoints and credentials are required, with no ambient AWS credential +lookup. The package is server-only and requires Node.js 20+. + +## Configure a provider + +```ts +import { S3Storage } from '@cleverbrush/storage-s3'; + +await using storage = new S3Storage({ + endpoint: process.env.STORAGE_ENDPOINT!, + region: process.env.STORAGE_REGION!, + bucket: process.env.STORAGE_BUCKET!, + credentials: { + accessKeyId: process.env.STORAGE_ACCESS_KEY_ID!, + secretAccessKey: process.env.STORAGE_SECRET_ACCESS_KEY! + }, + forcePathStyle: true, + keyPrefix: 'assets', + publicBaseUrl: 'https://assets.example.com' +}); + +await storage.put('images/logo.png', imageBytes, { contentType: 'image/png' }); +const url = storage.publicUrl('images/logo.png'); +// https://assets.example.com/assets/images/logo.png +// Leaving this scope automatically awaits storage.close(). +``` + +`await using` closes owned storage on both normal scope exit and exceptions, +waiting for active work and multipart cleanup. Keep application-wide instances +alive until shutdown; handlers borrowing injected storage must not dispose them. +Explicit `await storage.close()` is also supported and is idempotent. + +The public base URL is the public bucket/proxy root. The prefix is appended once, +so do not include the same prefix in both settings. No public URL is inferred +from the authenticated endpoint. Configuring it does not change bucket access. + +| Option | Default / purpose | +| --- | --- | +| `endpoint`, `region`, `bucket`, `credentials` | Required, explicit provider configuration. Credentials optionally include `sessionToken`. | +| `forcePathStyle` | `true`; set `false` for virtual-hosted bucket addressing. | +| `keyPrefix` | Empty; applied to every operation and public URL, omitted from returned keys. | +| `publicBaseUrl` | Unset; `publicUrl()` then returns `undefined`. | +| `partSize` | 5 MiB minimum and default. Increase for uploads that would exceed 10,000 parts. | +| `queueSize` | Four concurrent parts; multipart buffering is approximately part size × concurrency, plus stream buffers. | +| `requestTimeoutMs` | 30,000 milliseconds per network request, including cleanup; connections are limited to at most 10 seconds. | +| `requestChecksumCalculation` | `WHEN_REQUIRED`; optional `WHEN_SUPPORTED` for providers supporting SDK checksum extensions. | +| `responseChecksumValidation` | `WHEN_REQUIRED`; optional `WHEN_SUPPORTED`. | + +The SDK retries retryable requests up to three attempts. Streams are broken into +bounded, replayable parts instead of retrying a consumed stream as a whole. Failed +uploads attempt `AbortMultipartUpload`, including failures during completion. +Cleanup requires a reachable provider; configure an incomplete-upload lifecycle +rule as an operational backstop. `close()` waits for owned work and cleanup before +destroying the client. It is idempotent and also available via `await using`. + +Writes default to `application/octet-stream`. S3 user metadata accepts ASCII +values and names containing letters, digits, `_` and `-`; names are lowercased, +and duplicate names after lowercasing are rejected. Keys, including the prefix, +must fit S3's 1,024-byte limit. Copy uses native same-bucket `CopyObject`, preserving +headers and metadata. Provider copy-size limits apply (commonly 5 GiB); multipart +copy, version selection, ACLs, object tags and bucket provisioning are outside +this API. Deletion affects the current object; version-retention policy remains +provider configuration. + +## Provider examples and verification + +| Provider | Endpoint | Region | Addressing | +| --- | --- | --- | --- | +| Garage test service | `http://127.0.0.1:` | `garage` | Path style | +| Hetzner example | `https://fsn1.your-objectstorage.com` | `fsn1` | `forcePathStyle: false` | +| Other self-hosted services | Your S3 service URL | Server-configured region | Configurable | + +The same adapter serves every row. Garage v2.3.0 is exercised by CI, including +multipart abort cleanup. The Hetzner configuration follows its +[SDK examples](https://docs.hetzner.com/storage/object-storage/getting-started/using-libraries/) +and [supported operations](https://docs.hetzner.com/storage/object-storage/supported-actions/); +a live Hetzner account is not part of the default test suite. MinIO and other +implementations can run the same contract suite; they are not claimed as CI-tested. + +## Typed multipart upload handler + +The storage API accepts the bytes already provided by the Framework upload +contract. Storage SDK imports and configuration stay in backend modules. + +```ts +import { randomUUID } from 'node:crypto'; +import { any, array, object, string } from '@cleverbrush/schema'; +import { ActionResult, endpoint, file, type Handler } from '@cleverbrush/server'; +import type { ObjectStorage } from '@cleverbrush/storage'; + +const AssetStorageToken = any().hasType(); +const UploadAssets = endpoint.post('/assets') + .upload(object({ images: array(file()).minLength(1) })) + .inject({ storage: AssetStorageToken }) + .responses({ 201: array(object({ key: string(), url: string().optional() })) }); + +const upload: Handler = async ({ files }, { storage }) => { + const results = []; + for (const image of files.images) { + const key = `images/${randomUUID()}`; + await storage.put(key, image.buffer, { + contentType: image.mimeType, size: image.size + }); + results.push({ key, url: storage.publicUrl(key) }); + } + return ActionResult.created(results); +}; +``` + +Add application authorization, reference tracking and recovery policy around this +handler. Keep DTOs and upload contracts in browser-safe modules; import only +`ObjectStorage` types in backend code. Existing HTTP multipart parsing remains +buffered; stream-capable storage does not change that HTTP contract. + +## Stream a stored object to an HTTP response + +`pipeline` handles backpressure and destroys streams on failure/disconnect. Attach +cancellation before opening the object so disconnects also cancel pending reads: + +```ts +import { pipeline } from 'node:stream/promises'; +import { ActionResult } from '@cleverbrush/server'; + +// Inside an authorized backend handler with storage and a resolved object key: +return ActionResult.raw(async (_request, response) => { + const controller = new AbortController(); + const abort = () => controller.abort(); + response.once('close', abort); + try { + const object = await storage.get(key, { signal: controller.signal }); + try { + response.setHeader('content-type', object.contentType ?? 'application/octet-stream'); + response.setHeader('content-length', object.size); + await pipeline(object.body, response); + } finally { + object.body.destroy(); + } + } finally { + response.off('close', abort); + } +}); +``` + +## Run compatibility tests + +```sh +npm ci +npm run build +npm run test:storage:integration +``` + +The command starts a disposable official Garage container pinned by digest, runs +shared contract and real multipart tests, and removes the container and temporary +configuration afterward. Docker is required for the default local/CI fixture. + +To test an existing, explicitly designated test bucket, set +`STORAGE_TEST_ENDPOINT`, `STORAGE_TEST_REGION`, `STORAGE_TEST_BUCKET`, +`STORAGE_TEST_ACCESS_KEY_ID`, `STORAGE_TEST_SECRET_ACCESS_KEY`, and optionally +`STORAGE_TEST_FORCE_PATH_STYLE=false` and `STORAGE_TEST_PUBLIC_BASE_URL` before +running the same command. No container or bucket is created in this mode. Tests +write only beneath unique `framework-storage-tests//` prefixes and delete +their objects afterward. The credentials need bucket metadata access, object read/write/delete/copy and +multipart upload/list/abort permissions. No production provider is provisioned by +this package. diff --git a/libs/storage-s3/integration/storage.test.ts b/libs/storage-s3/integration/storage.test.ts new file mode 100644 index 00000000..345afb5d --- /dev/null +++ b/libs/storage-s3/integration/storage.test.ts @@ -0,0 +1,154 @@ +import { randomUUID } from 'node:crypto'; +import { PassThrough, Readable } from 'node:stream'; +import { + GetObjectCommand, + ListMultipartUploadsCommand, + S3Client +} from '@aws-sdk/client-s3'; +import { S3Storage, type S3StorageOptions } from '@cleverbrush/storage-s3'; +import { describe, expect, it } from 'vitest'; +import { storageContract } from '../../storage/testing/contract.js'; + +function configuration(): S3StorageOptions { + const env = process.env; + for (const name of [ + 'STORAGE_TEST_ENDPOINT', + 'STORAGE_TEST_REGION', + 'STORAGE_TEST_BUCKET', + 'STORAGE_TEST_ACCESS_KEY_ID', + 'STORAGE_TEST_SECRET_ACCESS_KEY' + ]) + if (!env[name]) + throw new Error( + `${name} is required; use npm run test:storage:integration` + ); + return { + endpoint: env.STORAGE_TEST_ENDPOINT!, + region: env.STORAGE_TEST_REGION!, + bucket: env.STORAGE_TEST_BUCKET!, + credentials: { + accessKeyId: env.STORAGE_TEST_ACCESS_KEY_ID!, + secretAccessKey: env.STORAGE_TEST_SECRET_ACCESS_KEY! + }, + forcePathStyle: env.STORAGE_TEST_FORCE_PATH_STYLE !== 'false', + keyPrefix: `framework-storage-tests/${randomUUID()}`, + publicBaseUrl: + env.STORAGE_TEST_PUBLIC_BASE_URL ?? + 'https://cdn.example.test/bucket', + requestTimeoutMs: 5000 + }; +} +function rawClient(config: S3StorageOptions) { + return new S3Client({ + ...config, + requestChecksumCalculation: 'WHEN_REQUIRED', + responseChecksumValidation: 'WHEN_REQUIRED' + }); +} + +storageContract(() => new S3Storage(configuration())); + +describe('S3-compatible transfers', () => { + it('round trips multipart streams and exposes only logical keys', async () => { + const config = configuration(); + const storage = new S3Storage(config); + const client = rawClient(config); + const key = 'large #?.bin'; + const bytes = Buffer.alloc(12 * 1024 * 1024 + 19); + for (let i = 0; i < bytes.length; i++) bytes[i] = i % 251; + try { + const input = Readable.from( + (function* () { + for (let i = 0; i < bytes.length; i += 32768) + yield bytes.subarray(i, i + 32768); + })() + ); + expect( + await storage.put(key, input, { + metadata: { origin: 'stream' } + }) + ).toMatchObject({ key, size: bytes.length }); + const response = await client.send( + new GetObjectCommand({ + Bucket: config.bucket, + Key: `${config.keyPrefix}/${key}` + }) + ); + expect( + Buffer.from(await response.Body!.transformToByteArray()).equals( + bytes + ) + ).toBe(true); + expect(storage.publicUrl(key)).toBe( + `${config.publicBaseUrl}/${config.keyPrefix}/large%20%23%3F.bin` + ); + } finally { + await storage.delete(key); + await storage.close(); + client.destroy(); + } + }); + it('aborts a multipart transfer without committing an object or leaving parts', async () => { + const config = configuration(); + const storage = new S3Storage({ ...config, queueSize: 1 }); + const client = rawClient(config); + const key = 'interrupted'; + const controller = new AbortController(); + const body = new PassThrough(); + const writing = storage.put(key, body, { signal: controller.signal }); + const rejected = expect(writing).rejects.toMatchObject({ + code: 'aborted' + }); + try { + body.write(Buffer.alloc(6 * 1024 * 1024)); + await expect + .poll( + async () => { + const listed = await client.send( + new ListMultipartUploadsCommand({ + Bucket: config.bucket, + Prefix: config.keyPrefix + }) + ); + return listed.Uploads?.length ?? 0; + }, + { timeout: 10000 } + ) + .toBe(1); + controller.abort(); + await rejected; + expect(body.destroyed).toBe(true); + expect(await storage.stat(key)).toBeUndefined(); + const listed = await client.send( + new ListMultipartUploadsCommand({ + Bucket: config.bucket, + Prefix: config.keyPrefix + }) + ); + expect(listed.Uploads ?? []).toHaveLength(0); + } finally { + controller.abort(); + await writing.catch(() => {}); + await storage.delete(key); + await storage.close(); + client.destroy(); + } + }); + it('maps rejected credentials without exposing secrets', async () => { + const config = configuration(); + const storage = new S3Storage({ + ...config, + credentials: { + accessKeyId: config.credentials.accessKeyId, + secretAccessKey: 'invalid-test-secret' + } + }); + try { + await expect(storage.stat('anything')).rejects.toMatchObject({ + code: 'access_denied' + }); + } finally { + await storage.close(); + } + }); +}); diff --git a/libs/storage-s3/package.json b/libs/storage-s3/package.json new file mode 100644 index 00000000..8699e4e5 --- /dev/null +++ b/libs/storage-s3/package.json @@ -0,0 +1,34 @@ +{ + "name": "@cleverbrush/storage-s3", + "version": "4.4.3", + "description": "S3-compatible object storage for self-hosted and hosted providers", + "type": "module", + "sideEffects": false, + "license": "BSD-3-Clause", + "engines": { + "node": ">=20" + }, + "files": [ + "dist" + ], + "main": "./dist/index.js", + "types": "./dist/index.d.ts", + "exports": { + ".": { + "types": "./dist/index.d.ts", + "import": "./dist/index.js" + } + }, + "devDependencies": { + "@types/node": "^25.4.0" + }, + "scripts": { + "build": "tsup && tsc --project tsconfig.build.json --emitDeclarationOnly", + "clean": "rm -rf dist tsconfig.tsbuildinfo" + }, + "dependencies": { + "@cleverbrush/storage": "^4.4.3", + "@aws-sdk/client-s3": "^3.1145.0", + "@aws-sdk/lib-storage": "^3.1145.0" + } +} diff --git a/libs/storage-s3/src/S3Storage.test-d.ts b/libs/storage-s3/src/S3Storage.test-d.ts new file mode 100644 index 00000000..53d7bd3b --- /dev/null +++ b/libs/storage-s3/src/S3Storage.test-d.ts @@ -0,0 +1,18 @@ +import type { ObjectStorage } from '@cleverbrush/storage'; +import { expectTypeOf, it } from 'vitest'; +import type { S3Storage, S3StorageOptions } from './index.js'; + +it('implements the provider-neutral contract with explicit provider configuration', () => { + expectTypeOf().toExtend(); + expectTypeOf().toEqualTypeOf<{ + accessKeyId: string; + secretAccessKey: string; + sessionToken?: string; + }>(); + // @ts-expect-error A custom endpoint is required, not inferred from AWS defaults. + const _missingEndpoint: S3StorageOptions = { + region: 'garage', + bucket: 'assets', + credentials: { accessKeyId: 'a', secretAccessKey: 'b' } + }; +}); diff --git a/libs/storage-s3/src/S3Storage.test.ts b/libs/storage-s3/src/S3Storage.test.ts new file mode 100644 index 00000000..098a70ad --- /dev/null +++ b/libs/storage-s3/src/S3Storage.test.ts @@ -0,0 +1,441 @@ +import { once } from 'node:events'; +import { PassThrough, Readable } from 'node:stream'; +import { + AbortMultipartUploadCommand, + CompleteMultipartUploadCommand, + CopyObjectCommand, + CreateMultipartUploadCommand, + GetObjectCommand, + HeadObjectCommand, + PutObjectCommand, + S3Client, + UploadPartCommand +} from '@aws-sdk/client-s3'; +import type { ObjectStorage } from '@cleverbrush/storage'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { S3Storage, type S3StorageOptions } from './S3Storage.js'; + +const options: S3StorageOptions = { + endpoint: 'http://localhost:3900', + region: 'garage', + bucket: 'assets', + credentials: { accessKeyId: 'test', secretAccessKey: 'secret' }, + keyPrefix: 'prefix/', + publicBaseUrl: 'https://cdn.example.test/bucket', + queueSize: 1 +}; +const stores: S3Storage[] = []; +function create(extra: Partial = {}) { + const storage = new S3Storage({ ...options, ...extra }); + stores.push(storage); + return storage; +} +afterEach(async () => { + await Promise.all(stores.splice(0).map(store => store.close())); + vi.restoreAllMocks(); +}); + +describe('S3 storage', () => { + it('uses explicit connection and compatibility options without changing public URLs', async () => { + let client: S3Client; + const send = vi + .spyOn(S3Client.prototype, 'send') + .mockImplementation(async function (this: S3Client) { + client = this; + return { ContentLength: 0 }; + } as any); + const storage = create({ forcePathStyle: false }); + await storage.stat('a b'); + expect(client!.config.forcePathStyle).toBe(false); + expect(await client!.config.region()).toBe('garage'); + expect(await client!.config.credentials()).toMatchObject( + options.credentials + ); + expect(await client!.config.requestChecksumCalculation()).toBe( + 'WHEN_REQUIRED' + ); + expect(await client!.config.responseChecksumValidation()).toBe( + 'WHEN_REQUIRED' + ); + expect(send.mock.calls[0][0].input).toMatchObject({ + Bucket: 'assets', + Key: 'prefix/a b' + }); + expect(storage.publicUrl('a b')).toBe( + 'https://cdn.example.test/bucket/prefix/a%20b' + ); + expect( + create({ publicBaseUrl: undefined }).publicUrl('a') + ).toBeUndefined(); + }); + it('keeps upload input bounded while the provider is slow', async () => { + let start!: () => void; + const ready = new Promise(resolve => { + start = resolve; + }); + let parts = 0; + vi.spyOn(S3Client.prototype, 'send').mockImplementation( + async (command: any, request: any) => { + if (command instanceof CreateMultipartUploadCommand) + return { UploadId: 'upload' }; + if (command instanceof UploadPartCommand) { + if (++parts === 2) start(); + await new Promise((_, reject) => + request.abortSignal.addEventListener( + 'abort', + () => reject({ name: 'AbortError' }), + { once: true } + ) + ); + } + return {}; + } + ); + let consumed = 0; + const chunk = Buffer.alloc(64 * 1024); + const input = Readable.from( + (function* () { + for (let i = 0; i < 1600; i++) { + consumed += chunk.length; + yield chunk; + } + })() + ); + const storage = create({ queueSize: 2 }); + const writing = storage.put('bounded', input); + const rejected = expect(writing).rejects.toMatchObject({ + code: 'aborted' + }); + await ready; + expect(consumed).toBeLessThan(4 * 5 * 1024 * 1024); + await storage.close(); + await rejected; + expect(input.destroyed).toBe(true); + }); + it.each([ + 'normal exit', + 'exception', + 'explicit close' + ])('await using closes unread downloads exactly once on %s', async exit => { + const source = new PassThrough(); + vi.spyOn(S3Client.prototype, 'send').mockResolvedValue({ + Body: source, + ContentLength: 20 + } as never); + const destroy = vi.spyOn(S3Client.prototype, 'destroy'); + const storage = create(); + const failure = new Error('scope failed'); + let body: Readable | undefined; + const scoped = (async () => { + await using owned: ObjectStorage = storage; + body = (await owned.get('unread')).body; + if (exit === 'explicit close') await owned.close(); + if (exit === 'exception') throw failure; + })(); + if (exit === 'exception') await expect(scoped).rejects.toBe(failure); + else await scoped; + expect(source.destroyed).toBe(true); + expect(body?.destroyed).toBe(true); + await expect(storage.stat('after-disposal')).rejects.toMatchObject({ + code: 'closed' + }); + await storage.close(); + expect(destroy).toHaveBeenCalledTimes(1); + }); + it('normalizes failed input streams and aborts their multipart state', async () => { + const source = new PassThrough(); + const send = vi + .spyOn(S3Client.prototype, 'send') + .mockImplementation(async (command: any) => { + if (command instanceof CreateMultipartUploadCommand) + return { UploadId: 'upload' }; + if (command instanceof UploadPartCommand) { + source.destroy(new Error('private source details')); + return { ETag: 'part' }; + } + return {}; + }); + const writing = create().put('source-error', source); + source.write(Buffer.alloc(6 * 1024 * 1024)); + await expect(writing).rejects.toMatchObject({ code: 'provider_error' }); + expect(send.mock.calls.at(-1)![0]).toBeInstanceOf( + AbortMultipartUploadCommand + ); + }); + it('owns source errors even when rejected before starting an upload', async () => { + const body = new Readable({ + read() {}, + destroy(_error, callback) { + callback(new Error('private pending open failure')); + } + }); + const closed = new Promise(resolve => + body.once('close', resolve) + ); + await expect( + create().put('key', body, { signal: AbortSignal.abort() }) + ).rejects.toMatchObject({ code: 'aborted' }); + await closed; + expect(body.destroyed).toBe(true); + expect(body.listenerCount('error')).toBe(0); + }); + it('releases malformed read responses even when destruction emits an error', async () => { + const source = new Readable({ + read() {}, + destroy(_error, callback) { + callback(new Error('private provider details')); + } + }); + const closed = new Promise(resolve => + source.once('close', resolve) + ); + vi.spyOn(S3Client.prototype, 'send').mockResolvedValue({ + Body: source + } as never); + await expect(create().get('invalid')).rejects.toMatchObject({ + code: 'provider_error' + }); + await closed; + expect(source.destroyed).toBe(true); + }); + it('copies prefixed reserved keys within one bucket, preserving metadata', async () => { + const send = vi + .spyOn(S3Client.prototype, 'send') + .mockImplementation(async (command: any) => + command instanceof HeadObjectCommand + ? { ContentLength: 7 } + : { CopyObjectResult: { ETag: 'opaque' } } + ); + expect(await create().copy('a #?.txt', 'b')).toEqual({ + key: 'b', + size: 7, + etag: 'opaque' + }); + expect(send.mock.calls[1][0]).toBeInstanceOf(CopyObjectCommand); + expect(send.mock.calls[1][0].input).toMatchObject({ + CopySource: 'assets/prefix/a%20%23%3F.txt', + Key: 'prefix/b', + MetadataDirective: 'COPY' + }); + }); + it('validates size and binary input without committing truncated data', async () => { + const send = vi.spyOn(S3Client.prototype, 'send'); + for (const [body, size] of [ + [Buffer.from('abc'), 2], + [Buffer.from('abc'), 4], + [Readable.from(['text']), undefined] + ] as const) { + await expect( + create().put('key', body, { size }) + ).rejects.toMatchObject({ code: 'invalid_argument' }); + } + expect(send).not.toHaveBeenCalled(); + }); + it('cleans up multipart uploads when completion fails', async () => { + const send = vi + .spyOn(S3Client.prototype, 'send') + .mockImplementation(async (command: any) => { + if (command instanceof CreateMultipartUploadCommand) + return { UploadId: 'upload' }; + if (command instanceof UploadPartCommand) + return { ETag: 'part' }; + if (command instanceof CompleteMultipartUploadCommand) + throw new Error('secret request'); + return {}; + }); + await expect( + create().put('key', Buffer.alloc(6 * 1024 * 1024)) + ).rejects.toMatchObject({ code: 'provider_error' }); + expect(send.mock.calls.at(-1)![0]).toBeInstanceOf( + AbortMultipartUploadCommand + ); + }); + it('aborts in-flight S3 requests and awaits multipart cleanup', async () => { + let started!: () => void; + const ready = new Promise(resolve => { + started = resolve; + }); + let cleaned = false; + vi.spyOn(S3Client.prototype, 'send').mockImplementation( + async (command: any, request: any) => { + if (command instanceof CreateMultipartUploadCommand) + return { UploadId: 'upload' }; + if (command instanceof UploadPartCommand) { + started(); + await new Promise((_, reject) => + request.abortSignal.addEventListener( + 'abort', + () => reject({ name: 'AbortError' }), + { once: true } + ) + ); + } + if (command instanceof AbortMultipartUploadCommand) { + expect(request.abortSignal.aborted).toBe(false); + cleaned = true; + } + return {}; + } + ); + const controller = new AbortController(); + const body = Readable.from([Buffer.alloc(6 * 1024 * 1024)]); + const write = create().put('key', body, { signal: controller.signal }); + const rejected = expect(write).rejects.toMatchObject({ + code: 'aborted' + }); + await ready; + controller.abort(); + await rejected; + expect(cleaned).toBe(true); + expect(body.destroyed).toBe(true); + }); + it('await using waits for multipart cleanup before leaving the scope', async () => { + let started!: () => void; + const ready = new Promise(resolve => { + started = resolve; + }); + let cleaning!: () => void; + const cleanupStarted = new Promise(resolve => { + cleaning = resolve; + }); + let finishCleanup!: () => void; + const cleanupFinished = new Promise(resolve => { + finishCleanup = resolve; + }); + vi.spyOn(S3Client.prototype, 'send').mockImplementation( + async (command: any, request: any) => { + if (command instanceof CreateMultipartUploadCommand) + return { UploadId: 'upload' }; + if (command instanceof UploadPartCommand) { + started(); + await new Promise((_, reject) => + request.abortSignal.addEventListener( + 'abort', + () => reject({ name: 'AbortError' }), + { once: true } + ) + ); + } + if (command instanceof AbortMultipartUploadCommand) { + expect(request.abortSignal.aborted).toBe(false); + cleaning(); + await cleanupFinished; + } + return {}; + } + ); + const destroy = vi.spyOn(S3Client.prototype, 'destroy'); + const storage = create(); + const body = Readable.from([Buffer.alloc(6 * 1024 * 1024)]); + let rejected!: Promise; + let exited = false; + const scoped = (async () => { + await using owned: ObjectStorage = storage; + rejected = expect(owned.put('key', body)).rejects.toMatchObject({ + code: 'aborted' + }); + await ready; + })().finally(() => { + exited = true; + }); + try { + await cleanupStarted; + await new Promise(resolve => setImmediate(resolve)); + expect(exited).toBe(false); + expect(destroy).not.toHaveBeenCalled(); + } finally { + finishCleanup(); + await scoped; + await rejected; + } + expect(exited).toBe(true); + expect(body.destroyed).toBe(true); + expect(destroy).toHaveBeenCalledTimes(1); + }); + it('cleans up a stalled input when closed and rejects subsequent operations', async () => { + const storage = create(); + const source = new PassThrough(); + const writing = storage.put('key', source); + const rejected = expect(writing).rejects.toMatchObject({ + code: 'aborted' + }); + await storage.close(); + await rejected; + expect(source.destroyed).toBe(true); + await expect(storage.stat('key')).rejects.toMatchObject({ + code: 'closed' + }); + await storage.close(); + }); + it('propagates source errors safely and releases abandoned download streams', async () => { + const source = new PassThrough(); + vi.spyOn(S3Client.prototype, 'send').mockResolvedValue({ + Body: source, + ContentLength: 20 + } as never); + const storage = create(); + const result = await storage.get('key'); + const error = once(result.body, 'error'); + source.destroy(new Error('secret provider details')); + expect((await error)[0]).toMatchObject({ code: 'provider_error' }); + await storage.close(); + expect(source.destroyed).toBe(true); + }); + it('cancels returned reads and closes the source when the consumer destroys the body', async () => { + const sources: PassThrough[] = []; + vi.spyOn(S3Client.prototype, 'send').mockImplementation( + async (command: any) => { + expect(command).toBeInstanceOf(GetObjectCommand); + const source = new PassThrough(); + sources.push(source); + return { Body: source, ContentLength: 20 }; + } + ); + const storage = create(); + const controller = new AbortController(); + const read = await storage.get('a', { signal: controller.signal }); + const failure = once(read.body, 'error'); + controller.abort(); + expect((await failure)[0]).toMatchObject({ code: 'aborted' }); + const other = await storage.get('b'); + other.body.destroy(); + await storage.close(); + expect(sources.every(source => source.destroyed)).toBe(true); + }); + it('never treats denied stat as missing or leaks provider errors', async () => { + vi.spyOn(S3Client.prototype, 'send').mockRejectedValue({ + $metadata: { httpStatusCode: 403 }, + request: { secret: 'secret' } + }); + await expect(create().stat('a')).rejects.toMatchObject({ + code: 'access_denied' + }); + }); + it('validates config and keys without making requests', async () => { + expect(() => + create({ endpoint: 'https://user:secret@example.test' }) + ).toThrow('Invalid storage argument'); + expect(() => create({ partSize: 1 })).toThrow( + 'Invalid storage argument' + ); + const send = vi.spyOn(S3Client.prototype, 'send'); + await expect( + create().put('../a', Buffer.from('x')) + ).rejects.toMatchObject({ code: 'invalid_argument' }); + expect(send).not.toHaveBeenCalled(); + }); + it('preserves empty writes and normalizes metadata names', async () => { + const send = vi + .spyOn(S3Client.prototype, 'send') + .mockResolvedValue({ ETag: 'opaque' } as never); + expect( + await create().put('empty', new Uint8Array(), { + metadata: { Owner: 'test' } + }) + ).toEqual({ key: 'empty', size: 0, etag: 'opaque' }); + expect(send.mock.calls[0][0]).toBeInstanceOf(PutObjectCommand); + expect(send.mock.calls[0][0].input).toMatchObject({ + Metadata: { owner: 'test' } + }); + }); +}); diff --git a/libs/storage-s3/src/S3Storage.ts b/libs/storage-s3/src/S3Storage.ts new file mode 100644 index 00000000..a852b05f --- /dev/null +++ b/libs/storage-s3/src/S3Storage.ts @@ -0,0 +1,468 @@ +import { PassThrough, Readable } from 'node:stream'; +import { + AbortMultipartUploadCommand, + CompleteMultipartUploadCommand, + CopyObjectCommand, + CreateMultipartUploadCommand, + DeleteObjectCommand, + GetObjectCommand, + HeadObjectCommand, + type HeadObjectCommandOutput, + S3Client +} from '@aws-sdk/client-s3'; +import { Upload } from '@aws-sdk/lib-storage'; +import { + type ObjectMetadata, + type ObjectStorage, + type PutOptions, + prefixedObjectKey, + publicObjectUrl, + type StorageBody, + StorageError, + type StorageOptions, + type StorageReadResult, + type StoredObject, + validateObjectKey, + validatePublicBaseUrl +} from '@cleverbrush/storage'; +import { storageError } from './errors.js'; + +/** Explicit credentials for the configured S3 service; no ambient AWS lookup. */ +export interface S3Credentials { + accessKeyId: string; + secretAccessKey: string; + sessionToken?: string; +} + +/** One adapter targets one bucket. Create separate instances for other providers. */ +export interface S3StorageOptions { + endpoint: string; + region: string; + bucket: string; + credentials: S3Credentials; + /** Defaults to true for self-hosted endpoints. */ + forcePathStyle?: boolean; + keyPrefix?: string; + /** Public bucket/proxy root; the encoded key prefix is appended automatically. */ + publicBaseUrl?: string; + /** Multipart buffer size, at least 5 MiB. Defaults to 5 MiB. */ + partSize?: number; + /** Number of concurrently buffered parts. Defaults to four. */ + queueSize?: number; + /** Per-request network timeout, including cleanup requests. Defaults to 30 seconds. */ + requestTimeoutMs?: number; + requestChecksumCalculation?: 'WHEN_REQUIRED' | 'WHEN_SUPPORTED'; + responseChecksumValidation?: 'WHEN_REQUIRED' | 'WHEN_SUPPORTED'; +} + +type Operation = { + controller: AbortController; + done: Promise; + release(): void; + retained: boolean; +}; + +/** S3 protocol adapter for custom endpoints, with bounded uploads and owned streams. */ +export class S3Storage implements ObjectStorage { + readonly #client: S3Client; + readonly #options: S3StorageOptions; + readonly #active = new Set(); + #closed = false; + #closing?: Promise; + + constructor(options: S3StorageOptions) { + validatePublicBaseUrl(options.endpoint); + if ( + !options.region || + !options.bucket || + /[/\\\s]/.test(options.bucket) || + !options.credentials?.accessKeyId || + !options.credentials.secretAccessKey + ) + throw new StorageError('invalid_argument'); + if (options.keyPrefix) validateObjectKey(options.keyPrefix); + if (options.publicBaseUrl !== undefined) + validatePublicBaseUrl(options.publicBaseUrl); + for (const [value, minimum] of [ + [options.partSize ?? 5 * 1024 * 1024, 5 * 1024 * 1024], + [options.queueSize ?? 4, 1], + [options.requestTimeoutMs ?? 30000, 1] + ]) { + if (!Number.isSafeInteger(value) || value < minimum) + throw new StorageError('invalid_argument'); + } + this.#options = { ...options, credentials: { ...options.credentials } }; + this.#client = new S3Client({ + endpoint: options.endpoint, + region: options.region, + credentials: { ...options.credentials }, + forcePathStyle: options.forcePathStyle ?? true, + requestChecksumCalculation: + options.requestChecksumCalculation ?? 'WHEN_REQUIRED', + responseChecksumValidation: + options.responseChecksumValidation ?? 'WHEN_REQUIRED', + maxAttempts: 3, + requestHandler: { + connectionTimeout: Math.min( + options.requestTimeoutMs ?? 30000, + 10000 + ), + requestTimeout: options.requestTimeoutMs ?? 30000, + throwOnRequestTimeout: true + } + }); + } + + #key(key: string): string { + const result = prefixedObjectKey(this.#options.keyPrefix ?? '', key); + if (Buffer.byteLength(result) > 1024) + throw new StorageError('invalid_argument'); + return result; + } + + async #run( + options: StorageOptions, + work: (operation: Operation) => Promise + ): Promise { + if (this.#closed) throw new StorageError('closed'); + if (options.signal?.aborted) throw new StorageError('aborted'); + const controller = new AbortController(); + const abort = () => controller.abort(); + let resolve!: () => void; + const operation: Operation = { + controller, + retained: false, + done: new Promise(done => { + resolve = done; + }), + release: () => { + options.signal?.removeEventListener('abort', abort); + this.#active.delete(operation); + resolve(); + } + }; + options.signal?.addEventListener('abort', abort, { once: true }); + this.#active.add(operation); + try { + return await work(operation); + } catch (error) { + throw storageError(error, controller.signal); + } finally { + if (!operation.retained) operation.release(); + } + } + + async put( + key: string, + body: StorageBody, + options: PutOptions = {} + ): Promise { + // Ownership starts before argument/abort checks: destroying a file stream + // can still report a pending open error even when no upload was started. + if (body instanceof Readable) { + const ignore = () => {}; + body.on('error', ignore); + body.once('close', () => body.off('error', ignore)); + } + try { + return await this.#run(options, async operation => { + const physical = this.#key(key); + if ( + !(body instanceof Uint8Array) && + !(body instanceof Readable) + ) + throw new StorageError('invalid_argument'); + if ( + options.size !== undefined && + (!Number.isSafeInteger(options.size) || options.size < 0) + ) + throw new StorageError('invalid_argument'); + const { signal } = operation.controller; + const metadata: Record = Object.create(null); + for (const [name, value] of Object.entries( + options.metadata ?? {} + )) { + const canonical = name.toLowerCase(); + if ( + !/^[a-z0-9][a-z0-9_-]*$/.test(canonical) || + typeof value !== 'string' || + /[^\x20-\x7e]/.test(value) || + Object.hasOwn(metadata, canonical) + ) + throw new StorageError('invalid_argument'); + metadata[canonical] = value; + } + let size = 0; + let uploadId: string | undefined; + const source = + body instanceof Readable ? body : Readable.from([body]); + // Own source errors even if cancellation happens before iteration starts. + const ignore = () => {}; + source.on('error', ignore); + const input = Readable.from( + (async function* () { + for await (const chunk of source) { + if (!(chunk instanceof Uint8Array)) + throw new StorageError('invalid_argument'); + size += chunk.byteLength; + if ( + !Number.isSafeInteger(size) || + (options.size !== undefined && + size > options.size) + ) + throw new StorageError('invalid_argument'); + yield chunk; + } + if (options.size !== undefined && size !== options.size) + throw new StorageError('invalid_argument'); + })(), + { objectMode: false } + ); + input.on('error', ignore); + const abort = () => { + source.destroy(new StorageError('aborted')); + input.destroy(new StorageError('aborted')); + }; + signal.addEventListener('abort', abort, { once: true }); + // Upload does not forward its AbortController to client.send(). + // Use an operation-local facade so every data request is cancellable, + // while cleanup remains possible after cancellation. Do not race + // Upload.abort(): await done() so workers and cleanup have settled. + const client = { + config: this.#client.config, + send: async (command: any) => { + const cleanup = + command instanceof AbortMultipartUploadCommand; + const result = (await this.#client.send(command, { + abortSignal: cleanup + ? AbortSignal.timeout( + this.#options.requestTimeoutMs ?? 30000 + ) + : signal + })) as any; + if (command instanceof CreateMultipartUploadCommand) + uploadId = result.UploadId; + if ( + cleanup || + command instanceof CompleteMultipartUploadCommand + ) + uploadId = undefined; + return result; + } + } as S3Client; + try { + const result = await new Upload({ + client, + params: { + Bucket: this.#options.bucket, + Key: physical, + Body: input, + ContentLength: options.size, + ContentType: + options.contentType ?? + 'application/octet-stream', + CacheControl: options.cacheControl, + ContentDisposition: options.contentDisposition, + Metadata: metadata + }, + partSize: this.#options.partSize ?? 5 * 1024 * 1024, + queueSize: this.#options.queueSize ?? 4, + leavePartsOnError: false + }).done(); + return { key, size, etag: result.ETag }; + } finally { + signal.removeEventListener('abort', abort); + source.destroy(); + input.destroy(); + // CompleteMultipartUpload failures are not cleaned up by Upload. + if (uploadId) { + await this.#client + .send( + new AbortMultipartUploadCommand({ + Bucket: this.#options.bucket, + Key: physical, + UploadId: uploadId + }), + { + abortSignal: AbortSignal.timeout( + this.#options.requestTimeoutMs ?? 30000 + ) + } + ) + .catch(() => {}); + } + } + }); + } finally { + if (body instanceof Readable) body.destroy(); + } + } + + #metadata(key: string, result: HeadObjectCommandOutput): ObjectMetadata { + if ( + !Number.isSafeInteger(result.ContentLength) || + result.ContentLength! < 0 + ) + throw new StorageError('provider_error'); + return { + key, + size: result.ContentLength!, + etag: result.ETag, + contentType: result.ContentType, + cacheControl: result.CacheControl, + contentDisposition: result.ContentDisposition, + metadata: { ...result.Metadata }, + lastModified: result.LastModified + }; + } + + async get( + key: string, + options: StorageOptions = {} + ): Promise { + return this.#run(options, async operation => { + const { signal } = operation.controller; + const result = await this.#client.send( + new GetObjectCommand({ + Bucket: this.#options.bucket, + Key: this.#key(key) + }), + { abortSignal: signal } + ); + const source = result.Body; + if (!(source instanceof Readable)) + throw new StorageError('provider_error'); + const ignore = () => {}; + source.on('error', ignore); + source.once('close', () => source.off('error', ignore)); + let metadata: ObjectMetadata; + try { + metadata = this.#metadata(key, result); + } catch (error) { + source.destroy(); + throw error; + } + const body = new PassThrough(); + const abort = () => body.destroy(new StorageError('aborted')); + body.on('error', () => {}); + source.on('error', error => + body.destroy(storageError(error, signal)) + ); + source.on('close', () => { + if (!source.readableEnded && !body.destroyed) + body.destroy(new StorageError('provider_error')); + }); + body.once('close', () => { + signal.removeEventListener('abort', abort); + source.destroy(); + operation.release(); + }); + operation.retained = true; + signal.addEventListener('abort', abort, { once: true }); + source.pipe(body); + if (signal.aborted) abort(); + return { ...metadata, body }; + }); + } + + async #head( + key: string, + signal: AbortSignal + ): Promise { + try { + const result = await this.#client.send( + new HeadObjectCommand({ + Bucket: this.#options.bucket, + Key: this.#key(key) + }), + { abortSignal: signal } + ); + return this.#metadata(key, result); + } catch (error) { + const normalized = storageError(error, signal); + if (normalized.code === 'not_found') return undefined; + throw normalized; + } + } + + stat( + key: string, + options: StorageOptions = {} + ): Promise { + return this.#run(options, operation => + this.#head(key, operation.controller.signal) + ); + } + + copy( + source: string, + destination: string, + options: StorageOptions = {} + ): Promise { + return this.#run(options, async operation => { + const from = this.#key(source); + const to = this.#key(destination); + const { signal } = operation.controller; + const metadata = await this.#head(source, signal); + if (!metadata) throw new StorageError('not_found'); + if (source === destination) return metadata; + const result = await this.#client.send( + new CopyObjectCommand({ + Bucket: this.#options.bucket, + Key: to, + CopySource: `${encodeURIComponent(this.#options.bucket)}/${from.split('/').map(encodeURIComponent).join('/')}`, + MetadataDirective: 'COPY' + }), + { abortSignal: signal } + ); + return { + key: destination, + size: metadata.size, + etag: result.CopyObjectResult?.ETag + }; + }); + } + + delete(key: string, options: StorageOptions = {}): Promise { + return this.#run(options, async operation => { + try { + await this.#client.send( + new DeleteObjectCommand({ + Bucket: this.#options.bucket, + Key: this.#key(key) + }), + { abortSignal: operation.controller.signal } + ); + } catch (error) { + const normalized = storageError( + error, + operation.controller.signal + ); + if (normalized.code !== 'not_found') throw normalized; + } + }); + } + + publicUrl(key: string): string | undefined { + const physical = this.#key(key); + return this.#options.publicBaseUrl === undefined + ? undefined + : publicObjectUrl(this.#options.publicBaseUrl, physical); + } + + close(): Promise { + if (!this.#closing) { + this.#closed = true; + const active = [...this.#active]; + for (const operation of active) operation.controller.abort(); + this.#closing = Promise.all( + active.map(operation => operation.done) + ).then(() => this.#client.destroy()); + } + return this.#closing; + } + + [Symbol.asyncDispose](): Promise { + return this.close(); + } +} diff --git a/libs/storage-s3/src/errors.test.ts b/libs/storage-s3/src/errors.test.ts new file mode 100644 index 00000000..06f989c2 --- /dev/null +++ b/libs/storage-s3/src/errors.test.ts @@ -0,0 +1,22 @@ +import { expect, it } from 'vitest'; +import { storageError } from './errors.js'; + +it.each([ + [{ $metadata: { httpStatusCode: 404 } }, 'not_found'], + [{ name: 'NoSuchKey' }, 'not_found'], + [{ $metadata: { httpStatusCode: 403 } }, 'access_denied'], + [{ name: 'SignatureDoesNotMatch' }, 'access_denied'], + [{ $metadata: { httpStatusCode: 503 } }, 'unavailable'], + [{ code: 'ECONNRESET' }, 'unavailable'], + [{ name: 'TimeoutError' }, 'unavailable'], + [{ name: 'AbortError' }, 'aborted'], + [ + { message: 'https://user:secret@example.test/?signature=secret' }, + 'provider_error' + ] +])('normalizes provider failures safely', (input, code) => { + const error = storageError(input); + expect(error.code).toBe(code); + expect(error.cause).toBeUndefined(); + expect(JSON.stringify(error) + error.stack).not.toContain('secret'); +}); diff --git a/libs/storage-s3/src/errors.ts b/libs/storage-s3/src/errors.ts new file mode 100644 index 00000000..0ebd944c --- /dev/null +++ b/libs/storage-s3/src/errors.ts @@ -0,0 +1,47 @@ +import { StorageError } from '@cleverbrush/storage'; + +/** Never retain a raw SDK/source error: it can contain signed URLs or credentials. */ +export function storageError( + error: unknown, + signal?: AbortSignal +): StorageError { + if (signal?.aborted) return new StorageError('aborted'); + if (error instanceof StorageError) return error; + const value = error as + | { + name?: string; + code?: string; + $metadata?: { httpStatusCode?: number }; + } + | undefined; + const status = value?.$metadata?.httpStatusCode; + if (value?.name === 'AbortError') return new StorageError('aborted'); + if ( + status === 404 || + value?.name === 'NoSuchKey' || + value?.name === 'NotFound' + ) + return new StorageError('not_found'); + if ( + status === 401 || + status === 403 || + [ + 'AccessDenied', + 'InvalidAccessKeyId', + 'SignatureDoesNotMatch' + ].includes(value?.name ?? '') + ) + return new StorageError('access_denied'); + if ( + (status !== undefined && status >= 500) || + status === 429 || + ['TimeoutError', 'RequestTimeout', 'SlowDown'].includes( + value?.name ?? '' + ) || + ['ECONNREFUSED', 'ECONNRESET', 'ETIMEDOUT', 'ENOTFOUND'].includes( + value?.code ?? '' + ) + ) + return new StorageError('unavailable'); + return new StorageError('provider_error'); +} diff --git a/libs/storage-s3/src/index.ts b/libs/storage-s3/src/index.ts new file mode 100644 index 00000000..981434a1 --- /dev/null +++ b/libs/storage-s3/src/index.ts @@ -0,0 +1,5 @@ +export { + type S3Credentials, + S3Storage, + type S3StorageOptions +} from './S3Storage.js'; diff --git a/libs/storage-s3/tsconfig.build.json b/libs/storage-s3/tsconfig.build.json new file mode 100644 index 00000000..41234fe7 --- /dev/null +++ b/libs/storage-s3/tsconfig.build.json @@ -0,0 +1,13 @@ +{ + "extends": "../../tsconfig.json", + "compilerOptions": { + "rootDir": "./src", + "outDir": "./dist", + "module": "ES2022", + "strict": true, + "skipLibCheck": true, + "types": ["node"] + }, + "include": ["src/**/*.ts"], + "exclude": ["src/**/*.test.ts", "src/**/*.test-d.ts"] +} diff --git a/libs/storage-s3/tsconfig.typecheck.json b/libs/storage-s3/tsconfig.typecheck.json new file mode 100644 index 00000000..8bdad1cf --- /dev/null +++ b/libs/storage-s3/tsconfig.typecheck.json @@ -0,0 +1,8 @@ +{ + "extends": "./tsconfig.build.json", + "compilerOptions": { + "noEmit": true + }, + "include": ["src/**/*.test-d.ts"], + "exclude": [] +} diff --git a/libs/storage-s3/tsup.config.ts b/libs/storage-s3/tsup.config.ts new file mode 100644 index 00000000..a79bfcde --- /dev/null +++ b/libs/storage-s3/tsup.config.ts @@ -0,0 +1,9 @@ +import { defineConfig } from 'tsup'; +export default defineConfig({ + entry: ['src/index.ts'], + format: ['esm'], + tsconfig: './tsconfig.build.json', + sourcemap: true, + clean: true, + target: 'es2022' +}); diff --git a/libs/storage-s3/vitest.config.mts b/libs/storage-s3/vitest.config.mts new file mode 100644 index 00000000..e18b86e9 --- /dev/null +++ b/libs/storage-s3/vitest.config.mts @@ -0,0 +1,11 @@ +import { defineConfig } from 'vitest/config'; +export default defineConfig({ + test: { + include: ['src/**/*.test.ts'], + typecheck: { + enabled: true, + include: ['src/**/*.test-d.ts'], + tsconfig: './tsconfig.typecheck.json' + } + } +}); diff --git a/libs/storage/README.md b/libs/storage/README.md new file mode 100644 index 00000000..8f8b5d08 --- /dev/null +++ b/libs/storage/README.md @@ -0,0 +1,79 @@ +# @cleverbrush/storage + +Provider-neutral object storage contracts and key/URL helpers for Node.js 20+. +The core has no runtime dependencies. Use an adapter such as +[`@cleverbrush/storage-s3`](../storage-s3) for self-hosted or hosted S3 services. + +## Contract + +| Method | Behavior | +| --- | --- | +| `put(key, bytesOrStream, options?)` | Replace an object; return logical key, actual byte size and optional ETag. | +| `get(key, { signal }?)` | Return metadata and a readable body; throw `not_found` when missing. | +| `stat(key, { signal }?)` | Return metadata, or `undefined` on a missing object; other errors throw. | +| `copy(source, destination, { signal }?)` | Copy within one storage instance, preserving metadata and replacing the destination. | +| `delete(key, { signal }?)` | Delete the current object; missing objects succeed. | +| `publicUrl(key)` | Return a stable public URL, or `undefined` without a configured public base URL. | +| `close()` | Abort active work, release resources and permanently close the instance. | +| `[Symbol.asyncDispose]()` | Await `close()` automatically when an `await using` scope ends. | + +`StorageBody` accepts `Uint8Array` (including Buffer) or a binary Node `Readable`. +`PutOptions` supports `size`, `contentType`, `cacheControl`, `contentDisposition`, +string `metadata` and `signal`. Supplying `size` asserts the exact byte count. +Metadata names are case-insensitive and returned in lowercase. `ObjectMetadata` +adds these headers and an optional `lastModified` date to the write receipt. +ETags are opaque; compare bytes or calculate an independent checksum when needed. + +## Ownership and errors + +The adapter owns a supplied write stream for that operation, consuming it once +and destroying it on failure or cancellation. A returned read stream belongs to +the caller: consume it, pipe it with `pipeline`, or destroy it. A read's abort +signal stays active until its body closes. Adapters must not buffer complete +streamed objects. + +Every `ObjectStorage` implements `AsyncDisposable`. Use `await using` when your +scope owns the instance; leaving that scope awaits `close()`, including when an +exception is thrown. Cleanup is asynchronous so active operations and multipart +cleanup finish before resources are released. Explicit `close()` remains +available and is safe to call more than once. + +`StorageError.code` is one of `not_found`, `access_denied`, `aborted`, +`invalid_argument`, `unavailable`, `closed` or `provider_error`. Errors omit raw +provider requests, signed URLs, credentials and source error messages. Permission +errors are not interpreted as absent objects. Cancellation can race a completed +write; object storage operations are not database transactions. + +## Keys and public URLs + +Keys are logical relative names, not filesystem paths or URLs. They are never +normalized or decoded. Empty keys, leading slashes, backslashes, control +characters and `.` / `..` path segments are rejected. Spaces, Unicode, literal +percent signs and repeated slashes are retained. `prefixedObjectKey` applies a +configured prefix once; `publicObjectUrl` encodes segments without encoding slash +separators. Public base URLs must be HTTP(S), without credentials, query or hash. + +A public URL is a mapping, not an access grant or an existence check. Configure +public bucket access or a proxy separately. Ownership checks, generated object +names, reference tracking, rendering and database consistency belong to the +application. Signed URLs and bucket management are not part of this contract. + +## Dependency injection + +Create application-owned schema tokens, just like other Framework services: + +```ts +import { any } from '@cleverbrush/schema'; +import type { ObjectStorage } from '@cleverbrush/storage'; + +export const AssetStorageToken = any().hasType(); +// services.addSingletonInstance(AssetStorageToken, storage); +// endpoint.inject({ storage: AssetStorageToken }); +``` + +Register multiple tokens/instances when assets use different providers or buckets. +Keep application-wide adapters alive until application shutdown, then close them +after stopping incoming work. Handlers borrowing injected storage must not dispose +it. Use `await using` only in the scope that owns the adapter's lifetime. +The reusable contract suite under `testing/` is a repository test utility and is +not included in the published package. diff --git a/libs/storage/package.json b/libs/storage/package.json new file mode 100644 index 00000000..17c05bd2 --- /dev/null +++ b/libs/storage/package.json @@ -0,0 +1,29 @@ +{ + "name": "@cleverbrush/storage", + "version": "4.4.3", + "description": "Provider-neutral object storage contracts for Node.js", + "type": "module", + "sideEffects": false, + "license": "BSD-3-Clause", + "engines": { + "node": ">=20" + }, + "files": [ + "dist" + ], + "main": "./dist/index.js", + "types": "./dist/index.d.ts", + "exports": { + ".": { + "types": "./dist/index.d.ts", + "import": "./dist/index.js" + } + }, + "devDependencies": { + "@types/node": "^25.4.0" + }, + "scripts": { + "build": "tsup && tsc --project tsconfig.build.json --emitDeclarationOnly", + "clean": "rm -rf dist tsconfig.tsbuildinfo" + } +} diff --git a/libs/storage/src/contract.test.ts b/libs/storage/src/contract.test.ts new file mode 100644 index 00000000..bb7aa793 --- /dev/null +++ b/libs/storage/src/contract.test.ts @@ -0,0 +1,88 @@ +import { Readable } from 'node:stream'; +import { storageContract } from '../testing/contract.js'; +import { + type ObjectMetadata, + type ObjectStorage, + publicObjectUrl, + StorageError, + type StorageOptions, + validateObjectKey +} from './index.js'; + +// Test-only reference adapter demonstrates that the contract needs no S3 SDK. +function memoryStorage(): ObjectStorage { + const values = new Map(); + let closed = false; + const check = (key: string, options: StorageOptions = {}) => { + if (closed) throw new StorageError('closed'); + if (options.signal?.aborted) throw new StorageError('aborted'); + validateObjectKey(key); + }; + return { + async put(key, body, options = {}) { + try { + check(key, options); + const chunks: Uint8Array[] = []; + if (body instanceof Uint8Array) chunks.push(body); + else + for await (const chunk of body) { + check(key, options); + chunks.push(chunk); + } + const bytes = Buffer.concat(chunks); + const { signal: _signal, ...headers } = options; + const info = { + ...headers, + key, + size: bytes.length, + metadata: Object.fromEntries( + Object.entries(options.metadata ?? {}).map(([k, v]) => [ + k.toLowerCase(), + v + ]) + ) + }; + values.set(key, { bytes, info }); + return { key, size: bytes.length }; + } finally { + if (body instanceof Readable) body.destroy(); + } + }, + async get(key, options) { + check(key, options); + const value = values.get(key); + if (!value) throw new StorageError('not_found'); + return { + ...structuredClone(value.info), + body: Readable.from([Buffer.from(value.bytes)]) + }; + }, + async stat(key, options) { + check(key, options); + return structuredClone(values.get(key)?.info); + }, + async copy(source, destination, options) { + check(source, options); + check(destination, options); + const value = values.get(source); + if (!value) throw new StorageError('not_found'); + const info = { ...structuredClone(value.info), key: destination }; + values.set(destination, { bytes: Buffer.from(value.bytes), info }); + return info; + }, + async delete(key, options) { + check(key, options); + values.delete(key); + }, + publicUrl(key) { + return publicObjectUrl('https://cdn.example.test', key); + }, + async close() { + closed = true; + }, + [Symbol.asyncDispose]() { + return this.close(); + } + }; +} +storageContract(memoryStorage); diff --git a/libs/storage/src/contracts.test-d.ts b/libs/storage/src/contracts.test-d.ts new file mode 100644 index 00000000..fc962c8e --- /dev/null +++ b/libs/storage/src/contracts.test-d.ts @@ -0,0 +1,21 @@ +import type { Readable } from 'node:stream'; +import { expectTypeOf, it } from 'vitest'; +import type { ObjectStorage, StorageBody } from './index.js'; + +declare const storage: ObjectStorage; + +it('supports await using through the provider-neutral type', async () => { + await using owned = storage; + expectTypeOf(owned).toEqualTypeOf(); + expectTypeOf().toExtend(); +}); + +it('accepts Node streams and bytes without provider types', () => { + expectTypeOf().toExtend(); + expectTypeOf().toExtend(); + expectTypeOf() + .returns.resolves.toHaveProperty('body') + .toEqualTypeOf(); + // @ts-expect-error Browser blobs are not server storage bodies. + const _blob: StorageBody = new Blob(); +}); diff --git a/libs/storage/src/contracts.ts b/libs/storage/src/contracts.ts new file mode 100644 index 00000000..04ba288d --- /dev/null +++ b/libs/storage/src/contracts.ts @@ -0,0 +1,75 @@ +import type { Readable } from 'node:stream'; + +/** Byte sources are consumed once; Node Buffers are Uint8Arrays. */ +export type StorageBody = Uint8Array | Readable; + +/** Cancellation remains active until a returned read stream is closed. */ +export interface StorageOptions { + signal?: AbortSignal; +} + +/** Portable HTTP and application metadata stored with an object. */ +export interface ObjectHeaders { + contentType?: string; + cacheControl?: string; + contentDisposition?: string; + /** Names are case-insensitive and returned in lowercase. */ + metadata?: Readonly>; +} + +/** An optional size must match the number of bytes consumed. */ +export interface PutOptions extends StorageOptions, ObjectHeaders { + size?: number; +} + +/** Receipt for a completed write or copy. ETags are not content checksums. */ +export interface StoredObject { + key: string; + size: number; + etag?: string; +} + +/** Metadata describes the current object, without opening a body stream. */ +export interface ObjectMetadata extends StoredObject, ObjectHeaders { + lastModified?: Date; +} + +/** The caller must consume or destroy body, even when inspecting only metadata. */ +export interface StorageReadResult extends ObjectMetadata { + body: Readable; +} + +/** + * Server-side, provider-neutral object storage. Keys are relative logical names; + * provider prefixes never appear in results. Operations on one key replace its + * current value; they do not provide database transactions or version history. + * Owners can use `await using` to close the instance when leaving its scope. + */ +export interface ObjectStorage extends AsyncDisposable { + put( + key: string, + body: StorageBody, + options?: PutOptions + ): Promise; + /** A missing object throws StorageError with code not_found. */ + get(key: string, options?: StorageOptions): Promise; + /** Only a confirmed missing object returns undefined; other failures throw. */ + stat( + key: string, + options?: StorageOptions + ): Promise; + /** Copy within this storage instance, replacing the destination and preserving metadata. */ + copy( + source: string, + destination: string, + options?: StorageOptions + ): Promise; + /** Deleting a missing object succeeds. */ + delete(key: string, options?: StorageOptions): Promise; + /** Stable URL only when a public base URL was explicitly configured. */ + publicUrl(key: string): string | undefined; + /** Cancel active work, await cleanup and permanently close this instance. Idempotent. */ + close(): Promise; + /** Delegate to close(), awaiting the same cleanup on scope exit. */ + [Symbol.asyncDispose](): Promise; +} diff --git a/libs/storage/src/errors.ts b/libs/storage/src/errors.ts new file mode 100644 index 00000000..7b383fe8 --- /dev/null +++ b/libs/storage/src/errors.ts @@ -0,0 +1,27 @@ +/** Portable failure classifications; no SDK request or credential data is exposed. */ +export type StorageErrorCode = + | 'not_found' + | 'access_denied' + | 'aborted' + | 'invalid_argument' + | 'unavailable' + | 'closed' + | 'provider_error'; + +const messages: Record = { + not_found: 'Storage object not found', + access_denied: 'Storage access denied', + aborted: 'Storage operation aborted', + invalid_argument: 'Invalid storage argument', + unavailable: 'Storage service unavailable', + closed: 'Storage is closed', + provider_error: 'Storage operation failed' +}; + +/** Safe to log: raw provider errors and request objects are deliberately omitted. */ +export class StorageError extends Error { + constructor(readonly code: StorageErrorCode) { + super(messages[code]); + this.name = 'StorageError'; + } +} diff --git a/libs/storage/src/index.ts b/libs/storage/src/index.ts new file mode 100644 index 00000000..b2445812 --- /dev/null +++ b/libs/storage/src/index.ts @@ -0,0 +1,8 @@ +export type * from './contracts.js'; +export { StorageError, type StorageErrorCode } from './errors.js'; +export { + prefixedObjectKey, + publicObjectUrl, + validateObjectKey, + validatePublicBaseUrl +} from './keys.js'; diff --git a/libs/storage/src/keys.test.ts b/libs/storage/src/keys.test.ts new file mode 100644 index 00000000..ded1a9af --- /dev/null +++ b/libs/storage/src/keys.test.ts @@ -0,0 +1,54 @@ +import { describe, expect, it } from 'vitest'; +import { StorageError } from './errors.js'; +import { + prefixedObjectKey, + publicObjectUrl, + validateObjectKey +} from './keys.js'; + +describe('object keys and public URLs', () => { + it('preserves opaque names and encodes reserved and Unicode characters once', () => { + expect(prefixedObjectKey('assets/', 'photo/雪 #?%2F.png')).toBe( + 'assets/photo/雪 #?%2F.png' + ); + expect( + publicObjectUrl( + 'https://cdn.example.test/bucket/', + 'assets/photo/雪 #?%2F.png' + ) + ).toBe( + 'https://cdn.example.test/bucket/assets/photo/%E9%9B%AA%20%23%3F%252F.png' + ); + expect(publicObjectUrl('https://cdn.example.test', 'a//b')).toBe( + 'https://cdn.example.test/a//b' + ); + expect(prefixedObjectKey('', 'a')).toBe('a'); + }); + it.each([ + '', + '/absolute', + '.', + '..', + 'a/../b', + 'a/./b', + 'a\\b', + 'a\0b', + '\ud800' + ])('rejects ambiguous key %j', key => { + expect(() => validateObjectKey(key)).toThrow(StorageError); + }); + it.each([ + 'ftp://example.test', + 'https://user:secret@example.test', + 'https://example.test?token=secret', + 'https://example.test/#x', + 'not a URL' + ])('rejects unsafe base URL %j', url => { + expect(() => publicObjectUrl(url, 'asset')).toThrow(StorageError); + try { + publicObjectUrl(url, 'asset'); + } catch (error) { + expect(String(error)).not.toContain('secret'); + } + }); +}); diff --git a/libs/storage/src/keys.ts b/libs/storage/src/keys.ts new file mode 100644 index 00000000..88a102e9 --- /dev/null +++ b/libs/storage/src/keys.ts @@ -0,0 +1,55 @@ +import { StorageError } from './errors.js'; + +/** Validate a logical key without normalizing or decoding its contents. */ +export function validateObjectKey(key: string): void { + if ( + typeof key !== 'string' || + !key || + key.startsWith('/') || + key.includes('\\') || + Array.from(key).some( + char => char.charCodeAt(0) < 32 || char.charCodeAt(0) === 127 + ) || + key.split('/').some(part => part === '.' || part === '..') + ) + throw new StorageError('invalid_argument'); + try { + encodeURIComponent(key); + } catch { + throw new StorageError('invalid_argument'); + } +} + +/** Join a configured prefix and a logical key exactly once. */ +export function prefixedObjectKey(prefix: string, key: string): string { + validateObjectKey(key); + if (!prefix) return key; + validateObjectKey(prefix); + return `${prefix.endsWith('/') ? prefix : `${prefix}/`}${key}`; +} + +/** Validate a public HTTP(S) base URL without exposing credentials in failures. */ +export function validatePublicBaseUrl(baseUrl: string): void { + let url: URL; + try { + url = new URL(baseUrl); + } catch { + throw new StorageError('invalid_argument'); + } + if ( + !['http:', 'https:'].includes(url.protocol) || + url.username || + url.password || + url.search || + url.hash + ) + throw new StorageError('invalid_argument'); +} + +/** Encode each object-key segment while retaining intentional slash separators. */ +export function publicObjectUrl(baseUrl: string, key: string): string { + validatePublicBaseUrl(baseUrl); + validateObjectKey(key); + const base = new URL(baseUrl).href.replace(/\/$/, ''); + return `${base}/${key.split('/').map(encodeURIComponent).join('/')}`; +} diff --git a/libs/storage/testing/contract.ts b/libs/storage/testing/contract.ts new file mode 100644 index 00000000..462e79a0 --- /dev/null +++ b/libs/storage/testing/contract.ts @@ -0,0 +1,132 @@ +import { Readable } from 'node:stream'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import type { ObjectStorage } from '../src/index.js'; + +/** Shared behavioral suite; intentionally not part of the published package. */ +export function storageContract(create: () => ObjectStorage) { + describe('object storage contract', () => { + let storage: ObjectStorage; + let keys: Set; + const key = (name: string) => { + keys.add(name); + return name; + }; + beforeEach(() => { + storage = create(); + keys = new Set(); + }); + afterEach(async () => { + try { + for (const name of keys) await storage.delete(name); + } finally { + await storage.close(); + } + }); + it('closes automatically at scope exit and remains safe to close again', async () => { + { + await using owned = storage; + expect(owned).toBe(storage); + } + await expect(storage.stat('after-disposal')).rejects.toMatchObject({ + code: 'closed' + }); + await storage.close(); + }); + it('round trips bytes and headers without treating ETags as checksums', async () => { + const name = key('folder/雪 #?%2F.bin'); + const bytes = Buffer.from([0, 1, 255, 42]); + expect( + await storage.put(name, bytes, { + contentType: 'application/x-test', + cacheControl: 'public, max-age=3600', + contentDisposition: 'inline', + metadata: { Owner: 'test' }, + size: bytes.length + }) + ).toMatchObject({ key: name, size: bytes.length }); + const expected = { + key: name, + size: bytes.length, + contentType: 'application/x-test', + cacheControl: 'public, max-age=3600', + contentDisposition: 'inline', + metadata: { owner: 'test' } + }; + expect(await storage.stat(name)).toMatchObject(expected); + const read = await storage.get(name); + expect(read).toMatchObject(expected); + const chunks = []; + for await (const chunk of read.body) chunks.push(chunk); + expect(Buffer.concat(chunks)).toEqual(bytes); + }); + it('supports empty bodies and replaces existing objects', async () => { + const name = key('replace'); + await storage.put(name, Buffer.from('old')); + await storage.put(name, new Uint8Array()); + expect(await storage.stat(name)).toMatchObject({ size: 0 }); + const read = await storage.get(name); + const chunks = []; + for await (const chunk of read.body) chunks.push(chunk); + expect(Buffer.concat(chunks)).toHaveLength(0); + }); + it('consumes unknown-length streams and preserves copied metadata', async () => { + const from = key('source #?.txt'); + const to = key('copy/result'); + const body = Readable.from([ + Buffer.from('first'), + Buffer.from('second') + ]); + await storage.put(from, body, { + contentType: 'text/plain', + metadata: { label: 'copy' } + }); + expect(body.destroyed).toBe(true); + await storage.put(to, Buffer.from('replaced')); + expect(await storage.copy(from, to)).toMatchObject({ + key: to, + size: 11 + }); + expect(await storage.stat(to)).toMatchObject({ + contentType: 'text/plain', + metadata: { label: 'copy' } + }); + const read = await storage.get(to); + const chunks = []; + for await (const chunk of read.body) chunks.push(chunk); + expect(Buffer.concat(chunks).toString()).toBe('firstsecond'); + }); + it('distinguishes missing reads from stat and idempotent deletion', async () => { + const name = key('missing'); + expect(await storage.stat(name)).toBeUndefined(); + await expect(storage.get(name)).rejects.toMatchObject({ + code: 'not_found' + }); + await expect( + storage.copy(name, key('copy-missing')) + ).rejects.toMatchObject({ code: 'not_found' }); + await storage.delete(name); + await storage.delete(name); + }); + it('rejects pre-aborted operations and closes supplied streams', async () => { + const signal = AbortSignal.abort(); + const name = key('aborted'); + const body = new Readable({ read() {} }); + await expect( + storage.put(name, body, { signal }) + ).rejects.toMatchObject({ code: 'aborted' }); + expect(body.destroyed).toBe(true); + await expect(storage.get(name, { signal })).rejects.toMatchObject({ + code: 'aborted' + }); + await expect(storage.stat(name, { signal })).rejects.toMatchObject({ + code: 'aborted' + }); + await expect( + storage.copy(name, name, { signal }) + ).rejects.toMatchObject({ code: 'aborted' }); + await expect( + storage.delete(name, { signal }) + ).rejects.toMatchObject({ code: 'aborted' }); + }); + }); +} diff --git a/libs/storage/tsconfig.build.json b/libs/storage/tsconfig.build.json new file mode 100644 index 00000000..41234fe7 --- /dev/null +++ b/libs/storage/tsconfig.build.json @@ -0,0 +1,13 @@ +{ + "extends": "../../tsconfig.json", + "compilerOptions": { + "rootDir": "./src", + "outDir": "./dist", + "module": "ES2022", + "strict": true, + "skipLibCheck": true, + "types": ["node"] + }, + "include": ["src/**/*.ts"], + "exclude": ["src/**/*.test.ts", "src/**/*.test-d.ts"] +} diff --git a/libs/storage/tsconfig.typecheck.json b/libs/storage/tsconfig.typecheck.json new file mode 100644 index 00000000..8bdad1cf --- /dev/null +++ b/libs/storage/tsconfig.typecheck.json @@ -0,0 +1,8 @@ +{ + "extends": "./tsconfig.build.json", + "compilerOptions": { + "noEmit": true + }, + "include": ["src/**/*.test-d.ts"], + "exclude": [] +} diff --git a/libs/storage/tsup.config.ts b/libs/storage/tsup.config.ts new file mode 100644 index 00000000..a79bfcde --- /dev/null +++ b/libs/storage/tsup.config.ts @@ -0,0 +1,9 @@ +import { defineConfig } from 'tsup'; +export default defineConfig({ + entry: ['src/index.ts'], + format: ['esm'], + tsconfig: './tsconfig.build.json', + sourcemap: true, + clean: true, + target: 'es2022' +}); diff --git a/libs/storage/vitest.config.mts b/libs/storage/vitest.config.mts new file mode 100644 index 00000000..e18b86e9 --- /dev/null +++ b/libs/storage/vitest.config.mts @@ -0,0 +1,11 @@ +import { defineConfig } from 'vitest/config'; +export default defineConfig({ + test: { + include: ['src/**/*.test.ts'], + typecheck: { + enabled: true, + include: ['src/**/*.test-d.ts'], + tsconfig: './tsconfig.typecheck.json' + } + } +}); diff --git a/package-lock.json b/package-lock.json index 7ddc70a1..4ad80dfa 100644 --- a/package-lock.json +++ b/package-lock.json @@ -546,6 +546,67 @@ "@cleverbrush/server": "^4.0.0" } }, + "libs/storage": { + "name": "@cleverbrush/storage", + "version": "4.4.3", + "license": "BSD-3-Clause", + "devDependencies": { + "@types/node": "^25.4.0" + }, + "engines": { + "node": ">=20" + } + }, + "libs/storage-s3": { + "name": "@cleverbrush/storage-s3", + "version": "4.4.3", + "license": "BSD-3-Clause", + "dependencies": { + "@aws-sdk/client-s3": "^3.1145.0", + "@aws-sdk/lib-storage": "^3.1145.0", + "@cleverbrush/storage": "^4.4.3" + }, + "devDependencies": { + "@types/node": "^25.4.0" + }, + "engines": { + "node": ">=20" + } + }, + "libs/storage-s3/node_modules/@types/node": { + "version": "25.9.9", + "resolved": "https://registry.npmjs.org/@types/node/-/node-25.9.9.tgz", + "integrity": "sha512-b4e2xxj/yMeT2hNlD6x7zNFiQK2Dlmkb9CQ6v3/1G2RezKN4YS2lcEv8lzGm1mhs4f2lM4yNAMI/zUDBZBBt7w==", + "dev": true, + "license": "MIT", + "dependencies": { + "undici-types": ">=7.24.0 <7.24.7" + } + }, + "libs/storage-s3/node_modules/undici-types": { + "version": "7.24.6", + "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-7.24.6.tgz", + "integrity": "sha512-WRNW+sJgj5OBN4/0JpHFqtqzhpbnV0GuB+OozA9gCL7a993SmU+1JBZCzLNxYsbMfIeDL+lTsphD5jN5N+n0zg==", + "dev": true, + "license": "MIT" + }, + "libs/storage/node_modules/@types/node": { + "version": "25.9.9", + "resolved": "https://registry.npmjs.org/@types/node/-/node-25.9.9.tgz", + "integrity": "sha512-b4e2xxj/yMeT2hNlD6x7zNFiQK2Dlmkb9CQ6v3/1G2RezKN4YS2lcEv8lzGm1mhs4f2lM4yNAMI/zUDBZBBt7w==", + "dev": true, + "license": "MIT", + "dependencies": { + "undici-types": ">=7.24.0 <7.24.7" + } + }, + "libs/storage/node_modules/undici-types": { + "version": "7.24.6", + "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-7.24.6.tgz", + "integrity": "sha512-WRNW+sJgj5OBN4/0JpHFqtqzhpbnV0GuB+OozA9gCL7a993SmU+1JBZCzLNxYsbMfIeDL+lTsphD5jN5N+n0zg==", + "dev": true, + "license": "MIT" + }, "node_modules/@asamuzakjp/css-color": { "version": "5.1.11", "resolved": "https://registry.npmjs.org/@asamuzakjp/css-color/-/css-color-5.1.11.tgz", @@ -597,6 +658,334 @@ "dev": true, "license": "MIT" }, + "node_modules/@aws-sdk/checksums": { + "version": "3.1001.1", + "resolved": "https://registry.npmjs.org/@aws-sdk/checksums/-/checksums-3.1001.1.tgz", + "integrity": "sha512-x12Q17KYlJAd3nKf8LV5LV0vt8sh8/6YfQLGPtrGnQf/tW4jqxPGq5GPpuVitpQYM3eUR4XB7CbxZf751NMbLw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/client-s3": { + "version": "3.1145.0", + "resolved": "https://registry.npmjs.org/@aws-sdk/client-s3/-/client-s3-3.1145.0.tgz", + "integrity": "sha512-MPKAQdx8qCZW1/ChDPILNehtAn/BMP4GNTnaH/xhGCWg3v50ADVq1K2WzxIrFfeANdWUyjwiY0CXDBhclsNWgQ==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/checksums": "^3.1001.1", + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/credential-provider-node": "^3.972.84", + "@aws-sdk/middleware-sdk-s3": "^3.972.77", + "@aws-sdk/signature-v4-multi-region": "^3.996.47", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/fetch-http-handler": "^5.8.0", + "@smithy/node-http-handler": "^4.12.1", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/core": { + "version": "3.978.1", + "resolved": "https://registry.npmjs.org/@aws-sdk/core/-/core-3.978.1.tgz", + "integrity": "sha512-LbY9aGsEiznDWmUc30Nwv3aIX/+dbwTx8KfS0yOC3NPYMO+O91e6jkT1azf34FwjOndq8/Q+RcVVZz5xnerwdg==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/types": "^3.974.6", + "@aws-sdk/xml-builder": "^3.972.41", + "@aws/lambda-invoke-store": "^0.3.0", + "@smithy/core": "^3.35.0", + "@smithy/signature-v4": "^5.7.3", + "@smithy/types": "^4.19.0", + "bowser": "^2.11.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-env": { + "version": "3.972.72", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-env/-/credential-provider-env-3.972.72.tgz", + "integrity": "sha512-xTKO/FWJPozTIXbozVnVGoNBhaGba8TBcx+KyUjRVeOlXE+dUc7GTR1cLvu0uTdIdmemzaFbqqCshXeZA1fZew==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-http": { + "version": "3.972.74", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-http/-/credential-provider-http-3.972.74.tgz", + "integrity": "sha512-u91E/hT8f4d1xy0Jl7VG4nVKJ3lxbrZkoBTeSVoJdWBiSEUMwMS/9+e0H/aJVQV//Lt5wuzP+E69v4aRSsNTmw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/fetch-http-handler": "^5.8.0", + "@smithy/node-http-handler": "^4.12.1", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-ini": { + "version": "3.973.17", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-ini/-/credential-provider-ini-3.973.17.tgz", + "integrity": "sha512-ged4KXdBkvIC81bLvNHHuQKdKak/VXhQTR1NWYTTqW0474nlmsxy9O/vlgTIohDDWH3xpBdtVMZRyjb+DnocDA==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/credential-provider-env": "^3.972.72", + "@aws-sdk/credential-provider-http": "^3.972.74", + "@aws-sdk/credential-provider-login": "^3.972.79", + "@aws-sdk/credential-provider-process": "^3.972.72", + "@aws-sdk/credential-provider-sso": "^3.973.16", + "@aws-sdk/credential-provider-web-identity": "^3.972.78", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/credential-provider-imds": "^4.5.2", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-login": { + "version": "3.972.79", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-login/-/credential-provider-login-3.972.79.tgz", + "integrity": "sha512-L+Z85anONJd8MaiuraO4wRxATCdEejBZ3K3eymzWI5JPXa9sOS9CkIm72PBKqXKX+Z9p9NGMX5AIMXm0LEflgw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-node": { + "version": "3.972.84", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-node/-/credential-provider-node-3.972.84.tgz", + "integrity": "sha512-oHt854odINVwzwsh+c5x69j0ajm4DbqqqVJ+O1ECsCIZeMDAbzFpXItaqP7UZstJj/ATdTk/KFSH0LaNAgV+kA==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/credential-provider-env": "^3.972.72", + "@aws-sdk/credential-provider-http": "^3.972.74", + "@aws-sdk/credential-provider-ini": "^3.973.17", + "@aws-sdk/credential-provider-process": "^3.972.72", + "@aws-sdk/credential-provider-sso": "^3.973.16", + "@aws-sdk/credential-provider-web-identity": "^3.972.78", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/credential-provider-imds": "^4.5.2", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-process": { + "version": "3.972.72", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-process/-/credential-provider-process-3.972.72.tgz", + "integrity": "sha512-rLIp2xbMjX/k9/od7APpqq1ZgXXnV0pOL1Th3ZsL8Wu0TRtBsDTVS8iPqcfRFcHakFxPvR04OSTv2ka2qOb/2A==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-sso": { + "version": "3.973.16", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-sso/-/credential-provider-sso-3.973.16.tgz", + "integrity": "sha512-IGihaJfFZYacJJr/odqILCoK7W/mvrZ7cuK7ECn3sAu4vLC6u0V8bS7mCGbdugJ8Aum2tnvqmx0F2MRFp2rn9g==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/token-providers": "3.1138.0", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/credential-provider-web-identity": { + "version": "3.972.78", + "resolved": "https://registry.npmjs.org/@aws-sdk/credential-provider-web-identity/-/credential-provider-web-identity-3.972.78.tgz", + "integrity": "sha512-/y9WvNtlcPBGLR0qc1a+9J/xtYZfVczvLUOuXaVWylzttH7ewsxwHtjmiJSolNrVSDorIxHGHMU61CbonRkmwA==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/lib-storage": { + "version": "3.1145.0", + "resolved": "https://registry.npmjs.org/@aws-sdk/lib-storage/-/lib-storage-3.1145.0.tgz", + "integrity": "sha512-sM6n9AYlOLpRt+qovDJWYZk2dDGS18dHaOB1QCJIibXKbEuM1lz/flf1WstHCbyrteX4ldgMRFLw0xxjbfSkAQ==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "buffer": "5.6.0", + "events": "3.3.0", + "stream-browserify": "3.0.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + }, + "peerDependencies": { + "@aws-sdk/client-s3": "^3.1145.0" + } + }, + "node_modules/@aws-sdk/middleware-sdk-s3": { + "version": "3.972.77", + "resolved": "https://registry.npmjs.org/@aws-sdk/middleware-sdk-s3/-/middleware-sdk-s3-3.972.77.tgz", + "integrity": "sha512-E7W2UOeUoc+lg3uIfR/dM7ZwusHwhBQrKMnlkRv4EXRR+C0YtV1pg25xC7GdZIhXH+NAMgZPCbE7o5to2cjFiw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/signature-v4-multi-region": "^3.996.47", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/nested-clients": { + "version": "3.997.46", + "resolved": "https://registry.npmjs.org/@aws-sdk/nested-clients/-/nested-clients-3.997.46.tgz", + "integrity": "sha512-oRxtBcka/JGHGs9l9p9IVajGoTP8vTPmoAzdHGy4Qcy9P5vPnDf6nhIeM/COQNY9k/OahImTRaLkHftoXvfcmQ==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/signature-v4-multi-region": "^3.996.47", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/fetch-http-handler": "^5.8.0", + "@smithy/node-http-handler": "^4.12.1", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/signature-v4-multi-region": { + "version": "3.996.47", + "resolved": "https://registry.npmjs.org/@aws-sdk/signature-v4-multi-region/-/signature-v4-multi-region-3.996.47.tgz", + "integrity": "sha512-Zk08macMvQTHzQJCLJVkOlviVoqwYMrpXv4lmLN7b7sAbiMoOK7Go0NYdR5UeF+MW8LIbRmwrNy9u/5VvX1U5g==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/types": "^3.974.6", + "@smithy/signature-v4": "^5.7.3", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/token-providers": { + "version": "3.1138.0", + "resolved": "https://registry.npmjs.org/@aws-sdk/token-providers/-/token-providers-3.1138.0.tgz", + "integrity": "sha512-GpyAr0DD63YOEmYFM6Df+gJuIgC92MMTiBK4FTKfxii5MJ9ge20epR7LyroulscYlG89J+ZB2ivFDPjvfQhzdw==", + "license": "Apache-2.0", + "dependencies": { + "@aws-sdk/core": "^3.978.1", + "@aws-sdk/nested-clients": "^3.997.46", + "@aws-sdk/types": "^3.974.6", + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/types": { + "version": "3.974.6", + "resolved": "https://registry.npmjs.org/@aws-sdk/types/-/types-3.974.6.tgz", + "integrity": "sha512-v/clNZzZnDxGyvpHMOGpJKVXFAExJzUNAAjaWGdcx8QAcXLGwTaOkw33p5SHAi0YAioK32xB3hWwOekRVfmfKg==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws-sdk/xml-builder": { + "version": "3.972.41", + "resolved": "https://registry.npmjs.org/@aws-sdk/xml-builder/-/xml-builder-3.972.41.tgz", + "integrity": "sha512-ctjVSyCMegrWfXlx6VqzSBFI6UqmQ5ZlnfMhdLIiWmhoH8UAQxSCP5N3OpG7X3k4LnS7ou74C4mt20+bfTW2aQ==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=20.0.0" + } + }, + "node_modules/@aws/lambda-invoke-store": { + "version": "0.3.0", + "resolved": "https://registry.npmjs.org/@aws/lambda-invoke-store/-/lambda-invoke-store-0.3.0.tgz", + "integrity": "sha512-sl4Bm6yiMNYrZKkqqDFWN0UfnWhlS8ivKxrYl+6t0gCLrqr8y3B2IqZZbFRkfaVVp7C/baApyh71P+LeE1A2sQ==", + "license": "Apache-2.0", + "engines": { + "node": ">=18.0.0" + } + }, "node_modules/@babel/code-frame": { "version": "7.29.0", "resolved": "https://registry.npmjs.org/@babel/code-frame/-/code-frame-7.29.0.tgz", @@ -1439,6 +1828,14 @@ "resolved": "libs/server-openapi", "link": true }, + "node_modules/@cleverbrush/storage": { + "resolved": "libs/storage", + "link": true + }, + "node_modules/@cleverbrush/storage-s3": { + "resolved": "libs/storage-s3", + "link": true + }, "node_modules/@cleverbrush/todo-backend": { "resolved": "demos/todo-backend", "link": true @@ -5617,6 +6014,87 @@ "integrity": "sha512-RNiOoTPkptFtSVzQevY/yWtZwf/RxyVnPy/OcA9HBM3MlGDnBEYL5B41H0MTn0Uec8Hi+2qUtTfG2WWZBmMejQ==", "license": "BSD-3-Clause" }, + "node_modules/@smithy/core": { + "version": "3.35.1", + "resolved": "https://registry.npmjs.org/@smithy/core/-/core-3.35.1.tgz", + "integrity": "sha512-i4YPS4B6ts7bjn7UwLnGjiZdprOvHvgGobFZsYK3GIY3E5hIqtj0rReU69BcTpGp+fvtraSNXeG1l+jtJvF55w==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/credential-provider-imds": { + "version": "4.5.2", + "resolved": "https://registry.npmjs.org/@smithy/credential-provider-imds/-/credential-provider-imds-4.5.2.tgz", + "integrity": "sha512-A9uSdn72ozbRUSit0eib0TW7nXuNPlaeM0zcGkJ+nE6tFcSDbnmtwoxbTCFBukVQcszDAyvsd7+rTduPTXpygg==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.33.2", + "@smithy/types": "^4.17.2", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/fetch-http-handler": { + "version": "5.8.0", + "resolved": "https://registry.npmjs.org/@smithy/fetch-http-handler/-/fetch-http-handler-5.8.0.tgz", + "integrity": "sha512-ycSJu3tFAQ4v04CBB0agqFMVsSQ1iG3yw+SpgxRqKfaURpQD4CZ8Wn0zPMmSnOuTpTh65Vz+EA0rMrw089wvkA==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.33.3", + "@smithy/types": "^4.18.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/node-http-handler": { + "version": "4.12.1", + "resolved": "https://registry.npmjs.org/@smithy/node-http-handler/-/node-http-handler-4.12.1.tgz", + "integrity": "sha512-ThMkboGeONWXAelq9FvGsuJC4rOi+qyC4/zhUF58xYpxUg5sQKx2VXZYJmtNjr4dSuBJ1HeJXETQILCz3wOHvw==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.33.3", + "@smithy/types": "^4.18.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/signature-v4": { + "version": "5.7.4", + "resolved": "https://registry.npmjs.org/@smithy/signature-v4/-/signature-v4-5.7.4.tgz", + "integrity": "sha512-tHy0K0VtqNd5Y7Y41h0a0Lhh0L1GzC08dTWg0F7vRJWFtTENg7IZikf3wQkanYIRdb7ngoIPMTmqgUi401fEeQ==", + "license": "Apache-2.0", + "dependencies": { + "@smithy/core": "^3.35.0", + "@smithy/types": "^4.19.0", + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, + "node_modules/@smithy/types": { + "version": "4.19.0", + "resolved": "https://registry.npmjs.org/@smithy/types/-/types-4.19.0.tgz", + "integrity": "sha512-r7jh49VJxGerfAcTQA6gXcKc+98zOp/tqRwzYjgOE+iSQsP6cEU1hq2QzbuipmP68QtYdY9wKEhiCQZIzHgZ4Q==", + "license": "Apache-2.0", + "dependencies": { + "tslib": "^2.6.2" + }, + "engines": { + "node": ">=18.0.0" + } + }, "node_modules/@standard-schema/spec": { "version": "1.1.0", "resolved": "https://registry.npmjs.org/@standard-schema/spec/-/spec-1.1.0.tgz", @@ -6488,6 +6966,12 @@ "node": "*" } }, + "node_modules/bowser": { + "version": "2.14.1", + "resolved": "https://registry.npmjs.org/bowser/-/bowser-2.14.1.tgz", + "integrity": "sha512-tzPjzCxygAKWFOJP011oxFHs57HzIhOEracIgAePE4pqB3LikALKnSzUyU4MGs9/iCEUuHlAJTjTc5M+u7YEGg==", + "license": "MIT" + }, "node_modules/brace-expansion": { "version": "5.0.5", "resolved": "https://registry.npmjs.org/brace-expansion/-/brace-expansion-5.0.5.tgz", @@ -6548,6 +7032,16 @@ "node": "^6 || ^7 || ^8 || ^9 || ^10 || ^11 || ^12 || >=13.7" } }, + "node_modules/buffer": { + "version": "5.6.0", + "resolved": "https://registry.npmjs.org/buffer/-/buffer-5.6.0.tgz", + "integrity": "sha512-/gDYp/UtU0eA1ys8bOs9J6a+E/KWIY+DZ+Q2WESNUA0jFRsJOc0SNUO6xJ5SGA1xueg3NL65W6s+NY5l9cunuw==", + "license": "MIT", + "dependencies": { + "base64-js": "^1.0.2", + "ieee754": "^1.1.4" + } + }, "node_modules/buffer-equal-constant-time": { "version": "1.0.1", "resolved": "https://registry.npmjs.org/buffer-equal-constant-time/-/buffer-equal-constant-time-1.0.1.tgz", @@ -7100,6 +7594,15 @@ "@types/estree": "^1.0.0" } }, + "node_modules/events": { + "version": "3.3.0", + "resolved": "https://registry.npmjs.org/events/-/events-3.3.0.tgz", + "integrity": "sha512-mQw+2fkQbALzQ7V0MY0IqdnXNOeTtP4r0lN9z7AAawCXgqea7bDii20AYrIBrFd/Hx0M2Ocz6S111CaFkUcb0Q==", + "license": "MIT", + "engines": { + "node": ">=0.8.x" + } + }, "node_modules/expect-type": { "version": "1.3.0", "resolved": "https://registry.npmjs.org/expect-type/-/expect-type-1.3.0.tgz", @@ -7480,6 +7983,26 @@ "url": "https://opencollective.com/express" } }, + "node_modules/ieee754": { + "version": "1.2.1", + "resolved": "https://registry.npmjs.org/ieee754/-/ieee754-1.2.1.tgz", + "integrity": "sha512-dcyqhDvX1C46lXZcVqCpK+FtMRQVdIMN6/Df5js2zouUsqG7I6sFxitIC+7KYK29KdXOLHdu9zL4sFnoVQnqaA==", + "funding": [ + { + "type": "github", + "url": "https://github.com/sponsors/feross" + }, + { + "type": "patreon", + "url": "https://www.patreon.com/feross" + }, + { + "type": "consulting", + "url": "https://feross.org/support" + } + ], + "license": "BSD-3-Clause" + }, "node_modules/ignore": { "version": "5.3.2", "resolved": "https://registry.npmjs.org/ignore/-/ignore-5.3.2.tgz", @@ -7505,6 +8028,12 @@ "node": ">=18" } }, + "node_modules/inherits": { + "version": "2.0.4", + "resolved": "https://registry.npmjs.org/inherits/-/inherits-2.0.4.tgz", + "integrity": "sha512-k/vGaX4/Yla3WzyMCvTQOXYeIHvqOKtnqBduzTHpzpQZzAskKMhZ2K+EnBiSM9zGSoIFeMpXKxa4dYeZIQqewQ==", + "license": "ISC" + }, "node_modules/interpret": { "version": "2.2.0", "resolved": "https://registry.npmjs.org/interpret/-/interpret-2.2.0.tgz", @@ -9154,6 +9683,20 @@ "js-yaml": "bin/js-yaml.js" } }, + "node_modules/readable-stream": { + "version": "3.6.2", + "resolved": "https://registry.npmjs.org/readable-stream/-/readable-stream-3.6.2.tgz", + "integrity": "sha512-9u/sniCrY3D5WdsERHzHE4G2YCXqoG5FTHUiCC4SIbr6XcLZBY05ya9EKjYek9O5xOAwjGq+1JdGBAS7Q9ScoA==", + "license": "MIT", + "dependencies": { + "inherits": "^2.0.3", + "string_decoder": "^1.1.1", + "util-deprecate": "^1.0.1" + }, + "engines": { + "node": ">= 6" + } + }, "node_modules/readdirp": { "version": "4.1.2", "resolved": "https://registry.npmjs.org/readdirp/-/readdirp-4.1.2.tgz", @@ -9583,6 +10126,25 @@ "dev": true, "license": "MIT" }, + "node_modules/stream-browserify": { + "version": "3.0.0", + "resolved": "https://registry.npmjs.org/stream-browserify/-/stream-browserify-3.0.0.tgz", + "integrity": "sha512-H73RAHsVBapbim0tU2JwwOiXUj+fikfiaoYAKHF3VJfA0pe2BCzkhAHBlLG6REzE+2WNZcxOXjK7lkso+9euLA==", + "license": "MIT", + "dependencies": { + "inherits": "~2.0.4", + "readable-stream": "^3.5.0" + } + }, + "node_modules/string_decoder": { + "version": "1.3.0", + "resolved": "https://registry.npmjs.org/string_decoder/-/string_decoder-1.3.0.tgz", + "integrity": "sha512-hkRX8U1WjJFd8LsDJ2yQ/wWWxaopEsABU1XfkM8A+j0+85JAGppt16cr1Whg6KIbb4okU6Mql6BOj+uup/wKeA==", + "license": "MIT", + "dependencies": { + "safe-buffer": "~5.2.0" + } + }, "node_modules/string-width": { "version": "4.2.3", "resolved": "https://registry.npmjs.org/string-width/-/string-width-4.2.3.tgz", @@ -10181,6 +10743,12 @@ "react": "^16.8.0 || ^17.0.0 || ^18.0.0 || ^19.0.0" } }, + "node_modules/util-deprecate": { + "version": "1.0.2", + "resolved": "https://registry.npmjs.org/util-deprecate/-/util-deprecate-1.0.2.tgz", + "integrity": "sha512-EPD5q1uXyFxJpCrLnCc1nHnq3gOa6DZBocAIiI2TaSCA7VCJ1UJDMagCzIkXNsUYfD1daK//LTEQ8xiIbrHtcw==", + "license": "MIT" + }, "node_modules/uuid": { "version": "9.0.1", "resolved": "https://registry.npmjs.org/uuid/-/uuid-9.0.1.tgz", diff --git a/package.json b/package.json index 8fc2e49b..cf5de23f 100644 --- a/package.json +++ b/package.json @@ -34,7 +34,8 @@ "test:e2e": "npm run test -w @cleverbrush/demo-e2e", "test:e2e:api": "npm run test:api -w @cleverbrush/demo-e2e", "test:e2e:ui": "npm run test:ui -w @cleverbrush/demo-e2e", - "test:e2e:reset": "RESET=1 npm run test -w @cleverbrush/demo-e2e" + "test:e2e:reset": "RESET=1 npm run test -w @cleverbrush/demo-e2e", + "test:storage:integration": "node scripts/test-storage.mjs" }, "workspaces": [ "./libs/*", diff --git a/scripts/test-storage.mjs b/scripts/test-storage.mjs new file mode 100644 index 00000000..3f4feaaf --- /dev/null +++ b/scripts/test-storage.mjs @@ -0,0 +1,88 @@ +import { randomBytes, randomUUID } from 'node:crypto'; +import { spawn, spawnSync } from 'node:child_process'; +import { mkdtemp, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { setTimeout as delay } from 'node:timers/promises'; +import { HeadBucketCommand, S3Client } from '@aws-sdk/client-s3'; + +const root = fileURLToPath(new URL('../', import.meta.url)); +const image = 'dxflrs/garage:v2.3.0@sha256:866bd13ed2038ba7e7190e840482bc27234c4afaf77be8cfa439ae088c1e4690'; +const env = { ...process.env }; +let container; +let directory; +let child; +const docker = args => { + const result = spawnSync('docker', args, { encoding: 'utf8' }); + if (result.error || result.status !== 0) throw new Error(`Docker failed: ${result.stderr || result.error}`); + return result.stdout.trim(); +}; +async function cleanup() { + if (container) { + spawnSync('docker', ['rm', '--force', '--volumes', container], { stdio: 'ignore' }); + container = undefined; + } + if (directory) await rm(directory, { recursive: true, force: true }); +} +for (const signal of ['SIGINT', 'SIGTERM']) { + process.once(signal, () => { + child?.kill(signal); + void cleanup().finally(() => process.exit(130)); + }); +} +try { + if (!env.STORAGE_TEST_ENDPOINT) { + directory = await mkdtemp(join(tmpdir(), 'framework-storage-')); + const config = join(directory, 'garage.toml'); + await writeFile(config, `metadata_dir = "/tmp/meta" +data_dir = "/tmp/data" +db_engine = "sqlite" +replication_factor = 1 +rpc_bind_addr = "[::]:3901" +rpc_public_addr = "127.0.0.1:3901" +rpc_secret = "${randomBytes(32).toString('hex')}" +[s3_api] +s3_region = "garage" +api_bind_addr = "[::]:3900" +`, { mode: 0o600 }); + Object.assign(env, { + STORAGE_TEST_REGION: 'garage', STORAGE_TEST_BUCKET: 'framework-storage-test', + STORAGE_TEST_ACCESS_KEY_ID: `GK${randomBytes(16).toString('hex')}`, + STORAGE_TEST_SECRET_ACCESS_KEY: randomBytes(32).toString('hex'), + STORAGE_TEST_FORCE_PATH_STYLE: 'true' + }); + container = `framework-storage-${randomUUID()}`; + docker(['run', '--detach', '--name', container, '-p', '127.0.0.1::3900', + '-v', `${config}:/etc/garage.toml:ro`, + '-e', `GARAGE_DEFAULT_ACCESS_KEY=${env.STORAGE_TEST_ACCESS_KEY_ID}`, + '-e', `GARAGE_DEFAULT_SECRET_KEY=${env.STORAGE_TEST_SECRET_ACCESS_KEY}`, + '-e', `GARAGE_DEFAULT_BUCKET=${env.STORAGE_TEST_BUCKET}`, + image, '/garage', 'server', '--single-node', '--default-bucket']); + env.STORAGE_TEST_ENDPOINT = `http://${docker(['port', container, '3900']).split('\n')[0]}`; + } + for (const name of ['STORAGE_TEST_REGION', 'STORAGE_TEST_BUCKET', 'STORAGE_TEST_ACCESS_KEY_ID', 'STORAGE_TEST_SECRET_ACCESS_KEY']) { + if (!env[name]) throw new Error(`${name} is required for the storage integration suite`); + } + const client = new S3Client({ + endpoint: env.STORAGE_TEST_ENDPOINT, region: env.STORAGE_TEST_REGION, + credentials: { accessKeyId: env.STORAGE_TEST_ACCESS_KEY_ID, secretAccessKey: env.STORAGE_TEST_SECRET_ACCESS_KEY }, + forcePathStyle: env.STORAGE_TEST_FORCE_PATH_STYLE !== 'false', maxAttempts: 1, + requestChecksumCalculation: 'WHEN_REQUIRED', responseChecksumValidation: 'WHEN_REQUIRED' + }); + try { + let ready = false; + for (let attempt = 0; attempt < 40; attempt++) { + try { + await client.send(new HeadBucketCommand({ Bucket: env.STORAGE_TEST_BUCKET }), { abortSignal: AbortSignal.timeout(1000) }); + ready = true; break; + } catch { await delay(250); } + } + if (!ready) throw new Error('S3 test bucket did not become ready; verify test endpoint, region and credentials'); + } finally { client.destroy(); } + child = spawn(process.execPath, [resolve(root, 'node_modules/vitest/vitest.mjs'), 'run', '--config', 'vitest.storage.config.mts'], { cwd: root, env, stdio: 'inherit' }); + process.exitCode = await new Promise((done, reject) => { child.once('error', reject); child.once('exit', code => done(code ?? 1)); }); +} catch (error) { + process.stderr.write(`${error.message}\n`); + process.exitCode = 1; +} finally { await cleanup(); } diff --git a/tsconfig.build.json b/tsconfig.build.json index 4adce918..432437a6 100644 --- a/tsconfig.build.json +++ b/tsconfig.build.json @@ -12,6 +12,8 @@ { "path": "libs/di" }, { "path": "libs/auth" }, { "path": "libs/server" }, - { "path": "libs/env/tsconfig.build.json" } + { "path": "libs/env/tsconfig.build.json" }, + { "path": "libs/storage/tsconfig.build.json" }, + { "path": "libs/storage-s3/tsconfig.build.json" } ] } diff --git a/vitest.storage.config.mts b/vitest.storage.config.mts new file mode 100644 index 00000000..82e69735 --- /dev/null +++ b/vitest.storage.config.mts @@ -0,0 +1,9 @@ +import { defineConfig } from 'vitest/config'; + +export default defineConfig({ + test: { + include: ['libs/storage-s3/integration/**/*.test.ts'], + testTimeout: 30000, + hookTimeout: 30000 + } +}); diff --git a/websites/docs/app/layout.tsx b/websites/docs/app/layout.tsx index 2af269c4..a2f57ecd 100644 --- a/websites/docs/app/layout.tsx +++ b/websites/docs/app/layout.tsx @@ -37,6 +37,7 @@ const NAV_ITEMS: NavItem[] = [ { href: '/react-form', label: 'React Form' }, { href: '/mapper', label: 'Mapper' }, { href: '/scheduler', label: 'Scheduler' }, + { href: '/storage', label: 'Storage' }, { href: '/log', label: 'Log' }, { href: '/otel', label: 'OpenTelemetry' } ] diff --git a/websites/docs/app/site.ts b/websites/docs/app/site.ts index 0e357374..9f7ea043 100644 --- a/websites/docs/app/site.ts +++ b/websites/docs/app/site.ts @@ -105,6 +105,12 @@ export const DOCS_ROUTES: RouteMetadata[] = [ description: 'Schema-driven object mapping with compile-time completeness and type-safe property selectors.' }, + { + path: '/storage', + title: 'Object Storage', + description: + 'Provider-neutral object storage and an S3-compatible adapter for self-hosted services and hosted providers.' + }, { path: '/scheduler', title: '@cleverbrush/scheduler', diff --git a/websites/docs/app/storage/page.tsx b/websites/docs/app/storage/page.tsx new file mode 100644 index 00000000..3ff1579f --- /dev/null +++ b/websites/docs/app/storage/page.tsx @@ -0,0 +1,87 @@ +import { docsMetadata } from '../site'; + +export const metadata = docsMetadata('/storage'); + +const setup = `import { S3Storage } from '@cleverbrush/storage-s3'; + +await using storage = new S3Storage({ + endpoint: process.env.STORAGE_ENDPOINT!, + region: process.env.STORAGE_REGION!, + bucket: process.env.STORAGE_BUCKET!, + credentials: { + accessKeyId: process.env.STORAGE_ACCESS_KEY_ID!, + secretAccessKey: process.env.STORAGE_SECRET_ACCESS_KEY! + }, + forcePathStyle: true, + keyPrefix: 'assets', + publicBaseUrl: 'https://assets.example.com' +}); + +await storage.put('images/logo.png', imageBytes, { + contentType: 'image/png' +}); +const url = storage.publicUrl('images/logo.png'); +// Leaving this scope automatically awaits storage.close().`; + +export default function StoragePage() { + return ( +
+

Object storage

+

+ @cleverbrush/storage defines a provider-neutral + Node.js contract. @cleverbrush/storage-s3{' '} + implements it for self-hosted and hosted S3-compatible services. + Change the endpoint, region, bucket and credentials to choose a + provider. +

+
+

Configure an adapter

+
+                    {setup}
+                
+

+ The public base URL is configured independently of the API + endpoint. It maps keys to stable public bucket or proxy + URLs; access policies remain deployment configuration. The + key prefix is appended once. Hetzner can use its regional + endpoint with virtual-hosted addressing by setting + forcePathStyle: false. +

+
+
+

Streaming and lifecycle

+

+ Use put, get, stat, copy and delete with an AbortSignal. + Writes accept buffers or binary Node streams and use bounded + multipart uploads. Consume or destroy every returned read + stream. Use await using for storage owned by + the current scope: normal exit and exceptions both await + active work and multipart cleanup before releasing + connections. Explicit close() is also supported + and is idempotent. Keep application-wide instances alive + until shutdown; handlers borrowing injected storage must not + dispose them. +

+

+ Metadata survives writes and copies. Errors have portable + codes and omit raw provider requests and credentials. Copy + stays within the configured bucket. ETags are opaque values, + not guaranteed content checksums. +

+

+ Garage is covered by the CI contract suite. Other endpoints, + including Hetzner, can run the same suite against a + designated test bucket. Production provider selection is + independent of the Garage test fixture. +

+ + Storage contract and DI + + {' · '} + + S3 configuration, uploads and streaming examples + +
+
+ ); +}