diff --git a/src/__tests__/upload.service.spec.ts b/src/__tests__/upload.service.spec.ts new file mode 100644 index 00000000..ef445e8b --- /dev/null +++ b/src/__tests__/upload.service.spec.ts @@ -0,0 +1,187 @@ +import axios from 'axios'; +import MockAdapter from 'axios-mock-adapter'; +import fs from 'fs-extra'; +import path from 'path'; +import UploadService from '../lib/api/services/upload.service'; +import { URLS } from '../lib/api/services/url'; +import { uninterceptedApiClient } from '../lib/api/ApiClient'; + +const assetsUploadData = require('./fixtures/assetsUploadData.json'); +const assetsCompleteUpload = require('./fixtures/assetsCompleteUpload.json'); + +describe('UploadService', () => { + let mock: MockAdapter; + let uploadMock: MockAdapter; + let tmpDir: string; + let filePath: string; + let previousApiRetryDelay: string; + let previousStorageRetryDelay: string; + let previousUploadConcurrency: string; + + beforeEach(() => { + previousApiRetryDelay = process.env.FDK_CLI_UPLOAD_API_RETRY_DELAY_MS; + previousStorageRetryDelay = + process.env.FDK_CLI_STORAGE_PUT_RETRY_DELAY_MS; + previousUploadConcurrency = process.env.FDK_CLI_UPLOAD_CONCURRENCY; + process.env.FDK_CLI_UPLOAD_API_RETRY_DELAY_MS = '0'; + process.env.FDK_CLI_STORAGE_PUT_RETRY_DELAY_MS = '0'; + mock = new MockAdapter(axios); + uploadMock = new MockAdapter(uninterceptedApiClient.axiosInstance); + tmpDir = fs.mkdtempSync(path.join(__dirname, 'upload-service-')); + filePath = path.join(tmpDir, 'themeBundle.css'); + fs.writeFileSync(filePath, 'body { color: #111; }'); + }); + + afterEach(() => { + mock.restore(); + uploadMock.restore(); + fs.removeSync(tmpDir); + restoreEnv('FDK_CLI_UPLOAD_API_RETRY_DELAY_MS', previousApiRetryDelay); + restoreEnv( + 'FDK_CLI_STORAGE_PUT_RETRY_DELAY_MS', + previousStorageRetryDelay, + ); + restoreEnv('FDK_CLI_UPLOAD_CONCURRENCY', previousUploadConcurrency); + }); + + it('retries transient upload/start 503 responses', async () => { + const namespace = 'application-theme-assets'; + const startUpload = buildStartUpload('https://upload.example.test/one.css'); + + mock.onPost(URLS.START_UPLOAD_FILE(namespace)).replyOnce( + 503, + 'upstream connect error or disconnect/reset before headers. reset reason: connection timeout', + { 'content-type': 'text/plain' }, + ); + mock.onPost(URLS.START_UPLOAD_FILE(namespace)).replyOnce(200, startUpload); + uploadMock.onPut(startUpload.upload.url).reply(200, ''); + mock.onPost(URLS.COMPLETE_UPLOAD_FILE(namespace)).reply( + 200, + assetsCompleteUpload, + ); + + const response = await UploadService.uploadFile(filePath, namespace); + + expect(response.complete).toEqual(assetsCompleteUpload); + expect( + mock.history.post.filter( + (request) => request.url === URLS.START_UPLOAD_FILE(namespace), + ), + ).toHaveLength(2); + }); + + it('retries transient Google Storage PUT DNS failures', async () => { + const namespace = 'application-theme-assets'; + const startUpload = buildStartUpload( + 'https://storage.googleapis.com/themeBundle.css', + ); + const dnsError: any = new Error( + 'getaddrinfo ENOTFOUND storage.googleapis.com', + ); + dnsError.code = 'ENOTFOUND'; + dnsError.request = {}; + dnsError.config = {}; + + mock.onPost(URLS.START_UPLOAD_FILE(namespace)).replyOnce(200, startUpload); + uploadMock.onPut(startUpload.upload.url).replyOnce(() => + Promise.reject(dnsError), + ); + uploadMock.onPut(startUpload.upload.url).replyOnce(200, ''); + mock.onPost(URLS.COMPLETE_UPLOAD_FILE(namespace)).reply( + 200, + assetsCompleteUpload, + ); + + const response = await UploadService.uploadFile(filePath, namespace); + + expect(response.complete).toEqual(assetsCompleteUpload); + expect(uploadMock.history.put).toHaveLength(2); + }); + + it('retries transient upload/complete 503 responses', async () => { + const namespace = 'application-theme-assets'; + const startUpload = buildStartUpload('https://upload.example.test/one.css'); + + mock.onPost(URLS.START_UPLOAD_FILE(namespace)).replyOnce(200, startUpload); + uploadMock.onPut(startUpload.upload.url).reply(200, ''); + mock.onPost(URLS.COMPLETE_UPLOAD_FILE(namespace)).replyOnce( + 503, + 'upstream connect error or disconnect/reset before headers. reset reason: connection timeout', + { 'content-type': 'text/plain' }, + ); + mock.onPost(URLS.COMPLETE_UPLOAD_FILE(namespace)).replyOnce( + 200, + assetsCompleteUpload, + ); + + const response = await UploadService.uploadFile(filePath, namespace); + + expect(response.complete).toEqual(assetsCompleteUpload); + expect( + mock.history.post.filter( + (request) => request.url === URLS.COMPLETE_UPLOAD_FILE(namespace), + ), + ).toHaveLength(2); + }); + + it('limits concurrent upload operations', async () => { + process.env.FDK_CLI_UPLOAD_CONCURRENCY = '1'; + const namespace = 'application-theme-assets'; + const secondFilePath = path.join(tmpDir, 'themeBundleTwo.css'); + fs.writeFileSync(secondFilePath, 'body { color: #222; }'); + let activePuts = 0; + let maxActivePuts = 0; + + mock.onPost(URLS.START_UPLOAD_FILE(namespace)).reply((config) => { + const data = + typeof config.data === 'string' + ? JSON.parse(config.data) + : config.data; + return [ + 200, + buildStartUpload( + `https://storage.googleapis.com/${data.file_name}`, + data.file_name, + ), + ]; + }); + uploadMock.onPut(/https:\/\/storage\.googleapis\.com\/.+/).reply( + async () => { + activePuts++; + maxActivePuts = Math.max(maxActivePuts, activePuts); + await new Promise((resolve) => setTimeout(resolve, 10)); + activePuts--; + return [200, '']; + }, + ); + mock.onPost(URLS.COMPLETE_UPLOAD_FILE(namespace)).reply( + 200, + assetsCompleteUpload, + ); + + await Promise.all([ + UploadService.uploadFile(filePath, namespace, 'one.css'), + UploadService.uploadFile(secondFilePath, namespace, 'two.css'), + ]); + + expect(uploadMock.history.put).toHaveLength(2); + expect(maxActivePuts).toBe(1); + }); +}); + +const buildStartUpload = (uploadUrl: string, fileName = 'themeBundle.css') => ({ + ...assetsUploadData, + file_name: fileName, + upload: { + ...assetsUploadData.upload, + url: uploadUrl, + }, +}); + +const restoreEnv = (key: string, value: string) => { + if (value === undefined) { + delete process.env[key]; + return; + } + process.env[key] = value; +}; diff --git a/src/helper/constants.ts b/src/helper/constants.ts index 46860e79..0c673516 100644 --- a/src/helper/constants.ts +++ b/src/helper/constants.ts @@ -14,6 +14,20 @@ export const ENVIRONMENT_COMMANDS = ['env']; export const AUTHENTICATION_COMMANDS = ['auth', 'login', 'logout']; export const EXTENSION_COMMANDS = ['init', 'get', 'set', 'pull-env']; export const MAX_RETRY = 5; +export const UPLOAD_API_MAX_ATTEMPTS = 3; +export const UPLOAD_API_RETRY_STATUS_CODES = [429, 502, 503, 504]; +export const UPLOAD_API_RETRY_DELAY_MS = 1000; +export const STORAGE_PUT_MAX_ATTEMPTS = 3; +export const STORAGE_PUT_RETRY_STATUS_CODES = [408, 429, 500, 502, 503, 504]; +export const STORAGE_PUT_RETRY_ERROR_CODES = [ + 'ENOTFOUND', + 'ETIMEDOUT', + 'ECONNRESET', + 'EPIPE', + 'ECONNABORTED', +]; +export const STORAGE_PUT_RETRY_DELAY_MS = 1000; +export const UPLOAD_CONCURRENCY = 16; export const THEME_TYPE = { vue2: 'vue2', react: 'react', @@ -111,4 +125,4 @@ export const PROJECT_REPOS = { [TEMPLATES['node-vue'].name]: TEMPLATES['node-vue'].repo, [TEMPLATES['node-react'].name]: TEMPLATES['node-react'].repo, [TEMPLATES['payment-node-react'].name]: TEMPLATES['payment-node-react'].repo -}; \ No newline at end of file +}; diff --git a/src/lib/api/services/upload.service.ts b/src/lib/api/services/upload.service.ts index 077532b2..b026d1ba 100644 --- a/src/lib/api/services/upload.service.ts +++ b/src/lib/api/services/upload.service.ts @@ -6,15 +6,162 @@ import fs from 'fs-extra'; import path from 'path'; import mime from 'mime'; import Spinner from '../../../helper/spinner'; +import Logger from '../../Logger'; +import { + STORAGE_PUT_MAX_ATTEMPTS, + STORAGE_PUT_RETRY_DELAY_MS, + STORAGE_PUT_RETRY_ERROR_CODES, + STORAGE_PUT_RETRY_STATUS_CODES, + UPLOAD_API_MAX_ATTEMPTS, + UPLOAD_API_RETRY_DELAY_MS, + UPLOAD_API_RETRY_STATUS_CODES, + UPLOAD_CONCURRENCY, +} from '../../../helper/constants'; + +let activeUploadSlots = 0; +const uploadQueue: Array<() => void> = []; + +const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); + +const getNumberFromEnv = ( + key: string, + fallback: number, + allowZero = false, +) => { + const value = Number(process.env[key]); + if (!Number.isFinite(value)) return fallback; + return allowZero ? (value >= 0 ? value : fallback) : value > 0 ? value : fallback; +}; + +const acquireUploadSlot = () => + new Promise((resolve) => { + const acquire = () => { + const limit = getNumberFromEnv( + 'FDK_CLI_UPLOAD_CONCURRENCY', + UPLOAD_CONCURRENCY, + ); + if (activeUploadSlots < limit) { + activeUploadSlots++; + resolve(); + return; + } + uploadQueue.push(acquire); + }; + acquire(); + }); + +const releaseUploadSlot = () => { + activeUploadSlots = Math.max(0, activeUploadSlots - 1); + const next = uploadQueue.shift(); + if (next) next(); +}; + +const getUploadApiRetryDelay = (attempt: number) => + Math.min( + getNumberFromEnv( + 'FDK_CLI_UPLOAD_API_RETRY_DELAY_MS', + UPLOAD_API_RETRY_DELAY_MS, + true, + ) * attempt, + 5000, + ); + +const getStoragePutRetryDelay = (attempt: number) => + Math.min( + getNumberFromEnv( + 'FDK_CLI_STORAGE_PUT_RETRY_DELAY_MS', + STORAGE_PUT_RETRY_DELAY_MS, + true, + ) * attempt, + 5000, + ); + +const isRetryableUploadApiError = (error) => + UPLOAD_API_RETRY_STATUS_CODES.indexOf(error?.response?.status) !== -1; + +const isRetryableStoragePutError = (error) => + STORAGE_PUT_RETRY_STATUS_CODES.indexOf(error?.response?.status) !== -1 || + STORAGE_PUT_RETRY_ERROR_CODES.indexOf(error?.code) !== -1; + +const redactSignedUrl = (url: string) => { + try { + const parsedUrl = new URL(url); + return `${parsedUrl.origin}${parsedUrl.pathname}`; + } catch (error) { + return url; + } +}; + +const postUploadApiWithRetry = async ({ + stage, + endpoint, + axiosOption, + debugId, + startData, + namespace, +}) => { + for (let attempt = 1; attempt <= UPLOAD_API_MAX_ATTEMPTS; attempt++) { + try { + return await ApiClient.post(endpoint, axiosOption); + } catch (error) { + if ( + attempt >= UPLOAD_API_MAX_ATTEMPTS || + !isRetryableUploadApiError(error) + ) { + throw error; + } + + const delayMs = getUploadApiRetryDelay(attempt); + Logger.debug( + `[uploadFile:${debugId}] ${stage} retry attempt=${attempt + 1}/${UPLOAD_API_MAX_ATTEMPTS} endpoint=${endpoint} file=${startData.file_name} namespace=${namespace} status=${error?.response?.status} message=${error?.message} delay_ms=${delayMs}`, + ); + await wait(delayMs); + } + } +}; + +const putStorageWithRetry = async ({ + s3Url, + data, + contentType, + debugId, + startData, + namespace, +}) => { + const endpoint = redactSignedUrl(s3Url); + for (let attempt = 1; attempt <= STORAGE_PUT_MAX_ATTEMPTS; attempt++) { + try { + return await uninterceptedApiClient.put(s3Url, { + data, + headers: { 'Content-Type': contentType }, + }); + } catch (error) { + if ( + attempt >= STORAGE_PUT_MAX_ATTEMPTS || + !isRetryableStoragePutError(error) + ) { + throw error; + } + + const delayMs = getStoragePutRetryDelay(attempt); + Logger.debug( + `[uploadFile:${debugId}] storage-put retry attempt=${attempt + 1}/${STORAGE_PUT_MAX_ATTEMPTS} endpoint=${endpoint} file=${startData.file_name} namespace=${namespace} code=${error?.code} status=${error?.response?.status} message=${error?.message} delay_ms=${delayMs}`, + ); + await wait(delayMs); + } + } +}; export default { uploadFile: async (filepath, namespace, file_name = null, mimeType = null) => { + await acquireUploadSlot(); let spinner if(process.env.DEBUG == 'fdk') { spinner = new Spinner(); } let textMessage; + const debugId = `${Date.now()}-${path.basename(filepath)}`; try { let stats = fs.statSync(filepath); textMessage = `Uploading file ${path.basename( @@ -38,19 +185,28 @@ export default { }, getCommonHeaderOptions(), ); - const res1 = await ApiClient.post( - URLS.START_UPLOAD_FILE(namespace), + const startEndpoint = URLS.START_UPLOAD_FILE(namespace); + const res1 = await postUploadApiWithRetry({ + stage: 'upload/start', + endpoint: startEndpoint, axiosOption, - ); + debugId, + startData, + namespace, + }); const startResponse = res1 ? res1.data : res1; let s3Url = startResponse.upload.url; // upload file to s3 // using uninterceptedApiClient to skip curl - const res2 = await uninterceptedApiClient.put(s3Url, { + const res2 = await putStorageWithRetry({ + s3Url, data: fs.readFileSync(filepath), - headers: { 'Content-Type': contentType }, + contentType, + debugId, + startData, + namespace, }); let uploadResponse = res2 ? res2.data : res2; @@ -65,10 +221,15 @@ export default { }, getCommonHeaderOptions(), ); - const res3 = await ApiClient.post( - URLS.COMPLETE_UPLOAD_FILE(namespace), + const completeEndpoint = URLS.COMPLETE_UPLOAD_FILE(namespace); + const res3 = await postUploadApiWithRetry({ + stage: 'upload/complete', + endpoint: completeEndpoint, axiosOption, - ); + debugId, + startData, + namespace, + }); let completeResponse = res3 ? res3.data : res3; spinner && spinner.succeed(textMessage); return { @@ -80,6 +241,7 @@ export default { spinner && spinner.fail(textMessage); throw error; } finally { + releaseUploadSlot(); spinner && spinner.stop(); } },