From b04fdc68dfbbd4597fb8d52a950fca0832ea75bd Mon Sep 17 00:00:00 2001 From: Hampus Date: Sat, 3 Oct 2026 02:11:34 +0200 Subject: [PATCH] fix(storage): fall back when cross-bucket copy is rejected (#3147) --- .../src/api/infrastructure/StorageService.ts | 91 +++++++- .../StorageServiceCrossBucketCopy.test.ts | 202 ++++++++++++++++++ 2 files changed, 286 insertions(+), 7 deletions(-) create mode 100644 fluxer_api/src/api/infrastructure/StorageServiceCrossBucketCopy.test.ts diff --git a/fluxer_api/src/api/infrastructure/StorageService.ts b/fluxer_api/src/api/infrastructure/StorageService.ts index 9710d5d2d..1613b1f67 100644 --- a/fluxer_api/src/api/infrastructure/StorageService.ts +++ b/fluxer_api/src/api/infrastructure/StorageService.ts @@ -108,6 +108,22 @@ function extractStreamFromGet(out: GetObjectCommandOutput): Readable { return wrapped; } +const REJECTED_SERVER_SIDE_COPY_ERRORS = new Set([ + 'NoSuchKey', + 'NotFound', + 'NotImplemented', + 'AccessDenied', + 'InvalidRequest', + 'MethodNotAllowed', +]); + +function isRejectedServerSideCopy(error: unknown): boolean { + return ( + error instanceof S3ServiceException && + (REJECTED_SERVER_SIDE_COPY_ERRORS.has(error.name) || error.$metadata?.httpStatusCode === 501) + ); +} + export class StorageService implements IStorageService { private readonly client: S3Client; private readonly presignClient: S3Client; @@ -473,15 +489,76 @@ export class StorageService implements IStorageService { if (isSameObject && !newContentType) { return; } - await this.client.send( - new CopyObjectCommand({ + try { + await this.client.send( + new CopyObjectCommand({ + Bucket: destinationBucket, + Key: destinationKey, + CopySource: `${encodeURIComponent(sourceBucket)}/${sourceKey.split('/').map(encodeURIComponent).join('/')}`, + ContentType: newContentType, + MetadataDirective: newContentType ? 'REPLACE' : undefined, + }), + ); + } catch (copyError) { + if (sourceBucket === destinationBucket || !isRejectedServerSideCopy(copyError)) { + throw copyError; + } + await this.copyObjectThroughApi( + {sourceBucket, sourceKey, destinationBucket, destinationKey, newContentType}, + copyError, + ); + } + } + + private async copyObjectThroughApi( + { + sourceBucket, + sourceKey, + destinationBucket, + destinationKey, + newContentType, + }: { + sourceBucket: string; + sourceKey: string; + destinationBucket: string; + destinationKey: string; + newContentType?: string; + }, + copyError: unknown, + ): Promise { + const source = await this.streamObject({bucket: sourceBucket, key: sourceKey}); + if (!source) { + throw copyError; + } + Logger.warn( + {sourceBucket, destinationBucket, error: copyError}, + 'Object storage rejected a cross-bucket copy, copying through the API instead', + ); + const upload = new Upload({ + client: this.client, + params: { Bucket: destinationBucket, Key: destinationKey, - CopySource: `${encodeURIComponent(sourceBucket)}/${sourceKey.split('/').map(encodeURIComponent).join('/')}`, - ContentType: newContentType, - MetadataDirective: newContentType ? 'REPLACE' : undefined, - }), - ); + Body: source.body, + ContentType: newContentType ?? source.contentType ?? undefined, + ...(newContentType + ? {} + : { + CacheControl: source.cacheControl ?? undefined, + ContentDisposition: source.contentDisposition ?? undefined, + Expires: source.expires ?? undefined, + }), + }, + partSize: STREAM_UPLOAD_PART_BYTES, + queueSize: STREAM_UPLOAD_CONCURRENCY, + leavePartsOnError: false, + }); + try { + await upload.done(); + } catch (error) { + source.body.destroy(); + throw error; + } } async copyObjectWithMetadataStripping({ diff --git a/fluxer_api/src/api/infrastructure/StorageServiceCrossBucketCopy.test.ts b/fluxer_api/src/api/infrastructure/StorageServiceCrossBucketCopy.test.ts new file mode 100644 index 000000000..740798af8 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/StorageServiceCrossBucketCopy.test.ts @@ -0,0 +1,202 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {StorageService} from '@app/api/infrastructure/StorageService'; +import {server} from '@app/api/test/msw/server'; +import {HttpResponse, http} from 'msw'; +import {beforeEach, describe, expect, it} from 'vitest'; + +const ENDPOINT = 'https://objects.ceph-rgw.test'; +const UPLOADS = 'fluxer-uploads'; +const CDN = 'fluxer-cdn'; + +interface StoredObject { + body: Uint8Array; + contentType: string; + cacheControl?: string; + contentDisposition?: string; +} + +const NO_SUCH_KEY = (bucket: string) => + new HttpResponse( + `NoSuchKey${bucket}tx0ceph`, + {status: 404, headers: {'Content-Type': 'application/xml'}}, + ); + +function cephWithoutCrossBucketCopy() { + const objects = new Map(); + const copies: Array<{source: string; destination: string}> = []; + const locate = (params: {bucket?: string | ReadonlyArray; key?: string | ReadonlyArray}) => { + const bucket = String(params.bucket); + const key = Array.isArray(params.key) ? params.key.join('/') : String(params.key); + return {bucket, key, id: `${bucket}/${key}`}; + }; + server.use( + http.put(`${ENDPOINT}/:bucket/*`, async ({request, params}) => { + const target = locate({bucket: params.bucket, key: params[0] as string}); + const copySource = request.headers.get('x-amz-copy-source'); + if (copySource) { + const decoded = decodeURIComponent(copySource.replace(/^\//u, '')); + const sourceBucket = decoded.split('/')[0]; + copies.push({source: decoded, destination: target.id}); + if (sourceBucket !== target.bucket) { + return NO_SUCH_KEY(target.bucket); + } + const source = objects.get(decoded); + if (!source) { + return NO_SUCH_KEY(target.bucket); + } + objects.set(target.id, { + ...source, + contentType: request.headers.get('content-type') ?? source.contentType, + }); + return new HttpResponse( + '"x"', + {status: 200, headers: {'Content-Type': 'application/xml'}}, + ); + } + const body = new Uint8Array(await request.arrayBuffer()); + objects.set(target.id, { + body, + contentType: request.headers.get('content-type') ?? 'application/octet-stream', + ...(request.headers.get('cache-control') ? {cacheControl: request.headers.get('cache-control')!} : {}), + ...(request.headers.get('content-disposition') + ? {contentDisposition: request.headers.get('content-disposition')!} + : {}), + }); + return new HttpResponse(null, {status: 200, headers: {ETag: '"x"'}}); + }), + http.get(`${ENDPOINT}/:bucket/*`, ({params}) => { + const target = locate({bucket: params.bucket, key: params[0] as string}); + const object = objects.get(target.id); + if (!object) { + return NO_SUCH_KEY(target.bucket); + } + return new HttpResponse(object.body, { + status: 200, + headers: { + 'Content-Type': object.contentType, + 'Content-Length': String(object.body.length), + ...(object.cacheControl ? {'Cache-Control': object.cacheControl} : {}), + ...(object.contentDisposition ? {'Content-Disposition': object.contentDisposition} : {}), + }, + }); + }), + http.head(`${ENDPOINT}/:bucket/*`, ({params}) => { + const target = locate({bucket: params.bucket, key: params[0] as string}); + const object = objects.get(target.id); + if (!object) { + return new HttpResponse(null, {status: 404}); + } + return new HttpResponse(null, { + status: 200, + headers: {'Content-Type': object.contentType, 'Content-Length': String(object.body.length)}, + }); + }), + ); + return {objects, copies}; +} + +function storage(): StorageService { + return new StorageService({ + endpoint: ENDPOINT, + forcePathStyle: true, + region: 'nbg1', + accessKeyId: 'TEST', + secretAccessKey: 'TEST', + }); +} + +const PDF = new TextEncoder().encode('%PDF-1.7\n1 0 obj << /Type /Catalog >> endobj\n%%EOF\n'); + +describe('StorageService copies on providers that reject cross-bucket CopyObject', () => { + let ceph: ReturnType; + + beforeEach(() => { + ceph = cephWithoutCrossBucketCopy(); + }); + + it('stores a non-media attachment in the CDN bucket', async () => { + ceph.objects.set(`${UPLOADS}/upload-1`, {body: PDF, contentType: 'application/octet-stream'}); + + await expect( + storage().copyObjectWithMetadataStripping({ + sourceBucket: UPLOADS, + sourceKey: 'upload-1', + destinationBucket: CDN, + destinationKey: 'attachments/1/2/file.pdf', + contentType: 'application/pdf', + }), + ).resolves.toBeNull(); + + const stored = ceph.objects.get(`${CDN}/attachments/1/2/file.pdf`); + expect(stored?.contentType).toBe('application/pdf'); + expect(Buffer.from(stored!.body).equals(Buffer.from(PDF))).toBe(true); + }); + + it('keeps the original file when media processing fails', async () => { + const brokenHeic = new Uint8Array(4096).fill(7); + ceph.objects.set(`${UPLOADS}/upload-2`, {body: brokenHeic, contentType: 'application/octet-stream'}); + + await expect( + storage().copyObjectWithMetadataStripping({ + sourceBucket: UPLOADS, + sourceKey: 'upload-2', + destinationBucket: CDN, + destinationKey: 'attachments/1/3/photo.heic', + contentType: 'image/heic', + }), + ).resolves.toBeNull(); + + const stored = ceph.objects.get(`${CDN}/attachments/1/3/photo.heic`); + expect(stored?.contentType).toBe('image/heic'); + expect(Buffer.from(stored!.body).equals(Buffer.from(brokenHeic))).toBe(true); + }); + + it('keeps the source headers when no new content type is given', async () => { + ceph.objects.set(`${CDN}/avatars/1/a.png`, { + body: PDF, + contentType: 'image/png', + cacheControl: 'public, max-age=31536000, immutable', + contentDisposition: 'inline', + }); + + await storage().copyObject({ + sourceBucket: CDN, + sourceKey: 'avatars/1/a.png', + destinationBucket: 'fluxer-reports', + destinationKey: 'evidence/a.png', + }); + + expect(ceph.objects.get('fluxer-reports/evidence/a.png')).toMatchObject({ + contentType: 'image/png', + cacheControl: 'public, max-age=31536000, immutable', + contentDisposition: 'inline', + }); + }); + + it('still fails when the source object does not exist', async () => { + await expect( + storage().copyObject({ + sourceBucket: UPLOADS, + sourceKey: 'missing', + destinationBucket: CDN, + destinationKey: 'attachments/1/4/missing.pdf', + newContentType: 'application/pdf', + }), + ).rejects.toMatchObject({name: 'NoSuchKey'}); + expect(ceph.objects.has(`${CDN}/attachments/1/4/missing.pdf`)).toBe(false); + }); + + it('does not retry a failed same-bucket copy through the API', async () => { + await expect( + storage().copyObject({ + sourceBucket: CDN, + sourceKey: 'missing', + destinationBucket: CDN, + destinationKey: 'other', + newContentType: 'image/png', + }), + ).rejects.toMatchObject({name: 'NoSuchKey'}); + expect(ceph.copies).toEqual([{source: `${CDN}/missing`, destination: `${CDN}/other`}]); + }); +});