diff --git a/.changeset/lemon-corners-hug.md b/.changeset/lemon-corners-hug.md new file mode 100644 index 0000000000..b41c09c928 --- /dev/null +++ b/.changeset/lemon-corners-hug.md @@ -0,0 +1,5 @@ +--- +'@backstage/plugin-techdocs-node': patch +--- + +There was an issue in the uploading of large size files to the AWS S3. We have modified the logic by adding retry along with multipart uploading functionality. diff --git a/plugins/techdocs-node/src/stages/publish/awsS3.test.ts b/plugins/techdocs-node/src/stages/publish/awsS3.test.ts index 6fb855d199..8be8c05994 100644 --- a/plugins/techdocs-node/src/stages/publish/awsS3.test.ts +++ b/plugins/techdocs-node/src/stages/publish/awsS3.test.ts @@ -22,6 +22,7 @@ import { ListObjectsV2Command, PutObjectCommand, S3Client, + S3ServiceException, UploadPartCommand, } from '@aws-sdk/client-s3'; import { Entity, DEFAULT_NAMESPACE } from '@backstage/catalog-model'; @@ -43,10 +44,13 @@ import { mockServices, } from '@backstage/backend-test-utils'; +jest.setTimeout(30_000); + const env = process.env; let s3Mock: AwsClientStub; -const mockDir = createMockDirectory(); +// Create a new MockDirectory for each test to avoid Windows file locking issues +let mockDir: ReturnType; function getMockCredentialProvider(): Promise { return Promise.resolve({ @@ -155,7 +159,7 @@ describe('AwsS3Publish', () => { build_timestamp: 612741599, }; - const directory = getEntityRootDir(entity); + let directory: string; const files = { 'index.html': '', @@ -176,7 +180,7 @@ describe('AwsS3Publish', () => { }, }; - beforeEach(() => { + beforeEach(async () => { process.env = { ...env }; process.env.AWS_REGION = 'us-west-2'; @@ -185,6 +189,12 @@ describe('AwsS3Publish', () => { getMockCredentialProvider(), ); + // Create a fresh mockdirectory for each test to avoid windows file locking + mockDir = createMockDirectory(); + // Calculate directory path with the new mockDir instance + directory = getEntityRootDir(entity); + + // Set up the test files mockDir.setContent({ [directory]: files, }); @@ -292,6 +302,161 @@ describe('AwsS3Publish', () => { }); }); + describe('retry mechanism', () => { + it('should retry with custom retry strategy', async () => { + const publisher = await createPublisherFromConfig(); + const customRetryStrategy = jest.fn(error => { + return error.name === 'NetworkingError'; + }); + + s3Mock + .on(ListObjectsV2Command) + .rejectsOnce( + new S3ServiceException({ + name: 'NetworkingError', + $fault: 'client', + $metadata: {}, + }), + ) + .resolvesOnce({ Contents: [] }); + + await (publisher as AwsS3Publish).retryOperation( + async () => { + const command = new ListObjectsV2Command({ Bucket: 'bucketName' }); + return (publisher as AwsS3Publish).storageClient.send(command); + }, + 'TestOperation', + 3, + customRetryStrategy, + ); + + expect(customRetryStrategy).toHaveBeenCalled(); + }); + + it('should use default retry strategy when no custom strategy provided', async () => { + const publisher = await createPublisherFromConfig(); + s3Mock + .on(ListObjectsV2Command) + .rejectsOnce( + new S3ServiceException({ + name: 'RequestTimeout', + $fault: 'client', + $metadata: {}, + }), + ) + .resolvesOnce({ Contents: [] }); + + await (publisher as AwsS3Publish).retryOperation( + async () => { + const command = new ListObjectsV2Command({ Bucket: 'bucketName' }); + return (publisher as AwsS3Publish).storageClient.send(command); + }, + 'TestOperation', + 3, + ); + }); + + it('should retry on server errors (5xx)', async () => { + const publisher = await createPublisherFromConfig(); + s3Mock + .on(ListObjectsV2Command) + .rejectsOnce( + new S3ServiceException({ + name: 'InternalError', + $fault: 'server', + $metadata: { httpStatusCode: 500 }, + }), + ) + .resolvesOnce({ Contents: [] }); + + await (publisher as AwsS3Publish).retryOperation( + async () => { + const command = new ListObjectsV2Command({ Bucket: 'bucketName' }); + return (publisher as AwsS3Publish).storageClient.send(command); + }, + 'TestOperation', + 3, + ); + }); + + it('should retry on specific 4xx errors that are transient', async () => { + const publisher = await createPublisherFromConfig(); + s3Mock + .on(ListObjectsV2Command) + .rejectsOnce( + new S3ServiceException({ + name: 'RequestTimeout', + $fault: 'client', + $metadata: { httpStatusCode: 408 }, + }), + ) + .resolvesOnce({ Contents: [] }); + + await (publisher as AwsS3Publish).retryOperation( + async () => { + const command = new ListObjectsV2Command({ Bucket: 'bucketName' }); + return (publisher as AwsS3Publish).storageClient.send(command); + }, + 'TestOperation', + 3, + ); + }); + + it('should use exact error code matching for transient errors', async () => { + const publisher = await createPublisherFromConfig(); + // Test that ConnectionError (exact match) is retried, but ConnectionErrorSomething (substring) is not + s3Mock + .on(ListObjectsV2Command) + .rejectsOnce( + new S3ServiceException({ + name: 'ConnectionError', + $fault: 'client', + $metadata: {}, + }), + ) + .resolvesOnce({ Contents: [] }); + + await (publisher as AwsS3Publish).retryOperation( + async () => { + const command = new ListObjectsV2Command({ Bucket: 'bucketName' }); + return (publisher as AwsS3Publish).storageClient.send(command); + }, + 'TestOperation', + 3, + ); + }); + + it('should apply exponential backoff with correct calculation', async () => { + const publisher = await createPublisherFromConfig(); + const startTime = Date.now(); + + s3Mock + .on(ListObjectsV2Command) + .rejectsOnce( + new S3ServiceException({ + name: 'SlowDown', + $fault: 'server', + $metadata: {}, + }), + ) + .resolvesOnce({ Contents: [] }); + + await (publisher as AwsS3Publish).retryOperation( + async () => { + const command = new ListObjectsV2Command({ Bucket: 'bucketName' }); + return (publisher as AwsS3Publish).storageClient.send(command); + }, + 'TestOperation', + 3, + ); + + const elapsedTime = Date.now() - startTime; + // First attempt fails, then backoff with baseDelay * 2^(attempt-1) = 1000 * 2^0 = 1000ms minimum + // Adding jitter (0-1000ms), so we expect at least 1000ms total + expect(elapsedTime).toBeGreaterThanOrEqual(900); + }); + }); + describe('getReadiness', () => { it('should validate correct config', async () => { const publisher = await createPublisherFromConfig(); @@ -405,7 +570,7 @@ describe('AwsS3Publish', () => { bucketName: bucketName, }); await publisher.publish({ entity, directory }); - expect(loggerInfoSpy).toHaveBeenLastCalledWith( + expect(loggerInfoSpy).toHaveBeenCalledWith( `Successfully deleted stale files for Entity ${entity.metadata.name}. Total number of files: 1`, ); }); @@ -549,7 +714,7 @@ describe('AwsS3Publish', () => { }); it('should return an error if the techdocs_metadata.json file cannot be read from stream', async () => { - s3Mock.on(GetObjectCommand).callsFake(_ => { + s3Mock.on(GetObjectCommand).callsFake(() => { return { Body: new ErrorReadable('No stream!'), }; @@ -688,7 +853,7 @@ describe('AwsS3Publish', () => { }); it('should return 404 if file cannot be read from stream', async () => { - s3Mock.on(GetObjectCommand).callsFake(_ => { + s3Mock.on(GetObjectCommand).callsFake(() => { return { Body: new ErrorReadable('No stream!'), }; diff --git a/plugins/techdocs-node/src/stages/publish/awsS3.ts b/plugins/techdocs-node/src/stages/publish/awsS3.ts index bd8b96ee14..157c2ffdeb 100644 --- a/plugins/techdocs-node/src/stages/publish/awsS3.ts +++ b/plugins/techdocs-node/src/stages/publish/awsS3.ts @@ -16,6 +16,10 @@ import { Entity, CompoundEntityRef } from '@backstage/catalog-model'; import { Config } from '@backstage/config'; import { assertError, ForwardedError } from '@backstage/errors'; + +// Maximum size in bytes for a single upload part (5MB) +const MAX_SINGLE_UPLOAD_BYTES = 5 * 1024 * 1024; + import { AwsCredentialsManager, DefaultAwsCredentialsManager, @@ -26,10 +30,12 @@ import { DeleteObjectCommand, HeadBucketCommand, HeadObjectCommand, + PutObjectCommand, PutObjectCommandInput, ListObjectsV2CommandOutput, ListObjectsV2Command, S3Client, + S3ServiceException, } from '@aws-sdk/client-s3'; import { fromTemporaryCredentials } from '@aws-sdk/credential-providers'; import { NodeHttpHandler } from '@smithy/node-http-handler'; @@ -78,12 +84,13 @@ const streamToBuffer = (stream: Readable): Promise => { }; export class AwsS3Publish implements PublisherBase { - private readonly storageClient: S3Client; + public readonly storageClient: S3Client; private readonly bucketName: string; private readonly legacyPathCasing: boolean; private readonly logger: LoggerService; private readonly bucketRootPath: string; private readonly sse?: 'aws:kms' | 'AES256'; + private readonly maxAttempts: number; constructor(options: { storageClient: S3Client; @@ -92,6 +99,7 @@ export class AwsS3Publish implements PublisherBase { logger: LoggerService; bucketRootPath: string; sse?: 'aws:kms' | 'AES256'; + maxAttempts: number; }) { this.storageClient = options.storageClient; this.bucketName = options.bucketName; @@ -99,6 +107,7 @@ export class AwsS3Publish implements PublisherBase { this.logger = options.logger; this.bucketRootPath = options.bucketRootPath; this.sse = options.sse; + this.maxAttempts = options.maxAttempts; } static async fromConfig( @@ -164,7 +173,7 @@ export class AwsS3Publish implements PublisherBase { 'techdocs.publisher.awsS3.s3ForcePathStyle', ); - // AWS MAX ATTEMPTS is an optional config. If missing, default value of 3 is used + // AWS MAX ATTEMPTS is an optional config. If missing, default value of 5 is used const maxAttempts = config.getOptionalNumber( 'techdocs.publisher.awsS3.maxAttempts', ); @@ -175,11 +184,16 @@ export class AwsS3Publish implements PublisherBase { ...(region && { region }), ...(endpoint && { endpoint }), ...(forcePathStyle && { forcePathStyle }), - ...(maxAttempts && { maxAttempts }), - ...(httpsProxy && { - requestHandler: new NodeHttpHandler({ + // Enhanced retry configuration for better reliability + maxAttempts: maxAttempts || 5, + retryMode: 'adaptive', + // Enhanced connection settings for large file uploads + requestHandler: new NodeHttpHandler({ + ...(httpsProxy && { httpsAgent: new HttpsProxyAgent({ proxy: httpsProxy }), }), + connectionTimeout: 60000, + socketTimeout: 120000, }), }); @@ -195,6 +209,7 @@ export class AwsS3Publish implements PublisherBase { legacyPathCasing, logger, sse, + maxAttempts: maxAttempts || 5, }); } @@ -250,6 +265,89 @@ export class AwsS3Publish implements PublisherBase { return explicitCredentials; } + /** + * Custom retry wrapper for S3 operations with detailed error handling. + */ + public async retryOperation( + operation: () => Promise, + operationName: string, + maxAttempts: number = 3, + shouldRetry: ( + error: S3ServiceException, + ) => boolean = this.defaultShouldRetry.bind(this), + ): Promise { + for (let attempt = 1; attempt < maxAttempts; attempt++) { + try { + return await operation(); + } catch (error) { + const e = error as S3ServiceException; + if (!shouldRetry(e)) { + this.logger.error(`${operationName} failed: ${e.message}`); + throw e; + } + + this.logger.warn(`${operationName} failed, retrying...`, { + attempt, + maxAttempts, + error: e.message, + errorCode: e.name, + httpStatusCode: e.$metadata?.httpStatusCode, + }); + + // Enhanced exponential backoff with jitter + const baseDelay = operationName.startsWith('Upload-') ? 2000 : 1000; + const backoffDelay = Math.min( + baseDelay * Math.pow(2, attempt - 1), + 30000, + ); + const jitter = Math.random() * 1000; + await new Promise(resolve => + setTimeout(resolve, backoffDelay + jitter), + ); + } + } + return await operation(); + } + + /** + * Determines if an S3 operation should be retried based on the error details. + */ + private defaultShouldRetry(error: S3ServiceException): boolean { + const httpStatusCode = error.$metadata?.httpStatusCode; + const errorCode = error.name; + + // Truly transient errors that should always be retried + const transientErrors = [ + 'NetworkingError', + 'TimeoutError', + 'ConnectionError', + 'RequestTimeout', + 'ServiceUnavailable', + 'SlowDown', + 'ThrottlingException', + ]; + + // Server errors are always considered transient + if (httpStatusCode && httpStatusCode >= 500) { + return true; + } + + // Specific 4xx errors that are known to be transient + if (httpStatusCode && httpStatusCode >= 400 && httpStatusCode < 500) { + const retriable4xxErrors = [ + 'RequestTimeout', + 'RequestTimeoutException', + 'PriorRequestNotComplete', + ]; + return retriable4xxErrors.includes(errorCode); + } + + // Check against known transient errors + return transientErrors.some( + retriableError => + errorCode === retriableError || error.message.includes(retriableError), + ); + } /** * Check if the defined bucket exists. Being able to connect means the configuration is good @@ -273,13 +371,15 @@ export class AwsS3Publish implements PublisherBase { 'explicitly defining credentials and region in techdocs.publisher.awsS3 in app config or ' + 'by using environment variables. Refer to https://backstage.io/docs/features/techdocs/using-cloud-storage', ); - this.logger.error(`from AWS client library`, error); + this.logger.error( + `from AWS client library`, + error instanceof Error ? error : new Error(String(error)), + ); return { isAvailable: false, }; } } - /** * Upload all the files from the generated `directory` to the S3 bucket. * Directory structure used in the bucket is - entityNamespace/entityKind/entityName/index.html @@ -293,6 +393,9 @@ export class AwsS3Publish implements PublisherBase { const bucketRootPath = this.bucketRootPath; const sse = this.sse; + // Track timing for performance monitoring + const publishStartTime = Date.now(); + // First, try to retrieve a list of all individual files currently existing let existingFiles: string[] = []; try { @@ -302,9 +405,20 @@ export class AwsS3Publish implements PublisherBase { useLegacyPathCasing, bucketRootPath, ); - existingFiles = await this.getAllObjectsFromBucket({ - prefix: remoteFolder, - }); + const response = await this.retryOperation( + async () => { + const listCommand = new ListObjectsV2Command({ + Bucket: this.bucketName, + Prefix: remoteFolder, + }); + return this.storageClient.send(listCommand); + }, + 'ListObjects', + this.maxAttempts, + ); + existingFiles = (response.Contents || []) + .map(f => f.Key || '') + .filter(f => !!f); } catch (e) { assertError(e); this.logger.error( @@ -323,27 +437,99 @@ export class AwsS3Publish implements PublisherBase { await bulkStorageOperation( async absoluteFilePath => { const relativeFilePath = path.relative(directory, absoluteFilePath); - const fileStream = fs.createReadStream(absoluteFilePath); - + const s3Key = getCloudPathForLocalPath( + entity, + relativeFilePath, + useLegacyPathCasing, + bucketRootPath, + ); + // Create params without the Body because the body must be the + // actual file contents (Buffer or Readable), not the path string. + // For multipart uploads we attach a Readable stream to avoid + // buffering large files in memory. For simple uploads we attach + // a Buffer read from disk. const params: PutObjectCommandInput = { Bucket: this.bucketName, - Key: getCloudPathForLocalPath( - entity, - relativeFilePath, - useLegacyPathCasing, - bucketRootPath, - ), - Body: fileStream, + Key: s3Key, ...(sse && { ServerSideEncryption: sse }), }; objects.push(params.Key!); + // Get file stats before upload + const stats = await fs.stat(absoluteFilePath); + const fileSizeInBytes = stats.size; - const upload = new Upload({ - client: this.storageClient, - params, - }); - return upload.done(); + // Check if this is a large file that requires multipart upload + if (fileSizeInBytes >= MAX_SINGLE_UPLOAD_BYTES) { + // Try multipart upload for large files + try { + // Create stream and Upload inside retry closure so stream is + // recreated on each retry attempt (streams are consumable once). + await this.retryOperation( + () => { + // Create a fresh stream on each attempt + const fileStream = fs.createReadStream(absoluteFilePath); + const uploadParams = { ...params, Body: fileStream }; + + const upload = new Upload({ + client: this.storageClient, + params: uploadParams, + partSize: MAX_SINGLE_UPLOAD_BYTES, + queueSize: 3, + leavePartsOnError: false, + }); + return upload.done(); + }, + `Upload-${params.Key}`, + this.maxAttempts, + ); + return; + } catch (multipartError) { + const s3Error = multipartError as any; + const errorName = s3Error?.name || 'Unknown'; + + // For specific multipart errors, attempt simple upload fallback + if (errorName === 'InvalidPart' || errorName === 'NoSuchUpload') { + this.logger.warn( + `Multipart upload failed for ${params.Key}, attempting simple upload fallback.`, + ); + } else { + // Non-recoverable multipart error, throw it + this.logger.error( + `Multipart upload failed for ${params.Key}: ${ + multipartError instanceof Error + ? multipartError.message + : String(multipartError) + }`, + ); + throw multipartError; + } + } + } + + // Use simple upload for small files or as fallback from multipart + try { + const fileContent = await fs.readFile(absoluteFilePath); + const putParams = { ...params, Body: fileContent }; + await this.retryOperation( + () => this.storageClient.send(new PutObjectCommand(putParams)), + `Upload-${params.Key}`, + this.maxAttempts, + ); + + if (fileSizeInBytes >= MAX_SINGLE_UPLOAD_BYTES) { + this.logger.info( + `Simple upload fallback succeeded for ${params.Key}`, + ); + } + } catch (error) { + this.logger.error( + `Upload failed for ${params.Key}: ${ + error instanceof Error ? error.message : String(error) + }`, + ); + throw error; + } }, absoluteFilesToUpload, { concurrencyLimit: 10 }, @@ -373,17 +559,21 @@ export class AwsS3Publish implements PublisherBase { await bulkStorageOperation( async relativeFilePath => { - return await this.storageClient.send( - new DeleteObjectCommand({ - Bucket: this.bucketName, - Key: relativeFilePath, - }), + return this.retryOperation( + async () => { + const deleteCommand = new DeleteObjectCommand({ + Bucket: this.bucketName, + Key: relativeFilePath, + }); + return this.storageClient.send(deleteCommand); + }, + 'DeleteObject', + this.maxAttempts, ); }, staleFiles, { concurrencyLimit: 10 }, ); - this.logger.info( `Successfully deleted stale files for Entity ${entity.metadata.name}. Total number of files: ${staleFiles.length}`, ); @@ -391,6 +581,13 @@ export class AwsS3Publish implements PublisherBase { const errorMessage = `Unable to delete file(s) from AWS S3. ${error}`; this.logger.error(errorMessage); } + const publishEndTime = Date.now(); + const publishDurationMs = publishEndTime - publishStartTime; + this.logger.info( + `Successfully published ${objects.length} files for ${ + entity.metadata.name + } in ${Math.round(publishDurationMs / 1000)}s`, + ); return { objects }; } @@ -413,11 +610,16 @@ export class AwsS3Publish implements PublisherBase { } try { - const resp = await this.storageClient.send( - new GetObjectCommand({ - Bucket: this.bucketName, - Key: `${entityRootDir}/techdocs_metadata.json`, - }), + const resp = await this.retryOperation( + async () => { + const getCommand = new GetObjectCommand({ + Bucket: this.bucketName, + Key: `${entityRootDir}/techdocs_metadata.json`, + }); + return this.storageClient.send(getCommand); + }, + 'GetTechDocsMetadata', + this.maxAttempts, ); const techdocsMetadataJson = await streamToBuffer( @@ -598,12 +800,18 @@ export class AwsS3Publish implements PublisherBase { let allObjects: ListObjectsV2CommandOutput; // Iterate through every file in the root of the publisher. do { - allObjects = await this.storageClient.send( - new ListObjectsV2Command({ - Bucket: this.bucketName, - ContinuationToken: nextContinuation, - ...(prefix ? { Prefix: prefix } : {}), - }), + const currentToken = nextContinuation; + allObjects = await this.retryOperation( + async () => { + const listCommand = new ListObjectsV2Command({ + Bucket: this.bucketName, + ContinuationToken: currentToken, + ...(prefix ? { Prefix: prefix } : {}), + }); + return this.storageClient.send(listCommand); + }, + 'GetAllObjects', + this.maxAttempts, ); objects.push( ...(allObjects.Contents || []).map(f => f.Key || '').filter(f => !!f),