diff --git a/src/app/network/UploadFolderManager.test.ts b/src/app/network/UploadFolderManager.test.ts index e184c3072f..120c2eceb8 100644 --- a/src/app/network/UploadFolderManager.test.ts +++ b/src/app/network/UploadFolderManager.test.ts @@ -12,6 +12,12 @@ import { import { FilesExceedsSizeLimitError } from 'app/drive/services/file.service/upload.errors'; import { uploadItemsParallelThunk } from 'app/store/slices/storage/storage.thunks/uploadItemsThunk'; import { deleteItemsThunk } from '../store/slices/storage/storage.thunks/deleteItemsThunk'; +import { uploadFileWithManager } from './UploadManager'; +import { PersistUploadRepository } from 'app/repositories/DatabaseUploadRepository'; + +vi.mock('app/drive/services/file.service/uploadFile', () => ({ + default: vi.fn(async (file: { name: string }) => ({ uuid: `${file.name}-uuid`, name: file.name })), +})); vi.mock('app/drive/services/new-storage.service', () => ({ default: { @@ -183,6 +189,76 @@ describe('uploadFoldersWithManager', () => { expect(deleteItemsThunk).not.toHaveBeenCalled(); }); + test('When reading the upload state of its files fails, then the folder upload finishes with failed files instead of hanging', async () => { + const mockFolder = buildFolderData({ name: 'MyFolder', plain_name: 'MyFolder' }); + const taskId = 'task-id'; + const failingUploadRepository: PersistUploadRepository = { + setUploadState: vi.fn(async () => undefined), + getUploadState: vi.fn(() => Promise.reject(new Error('IDB failure'))), + removeUploadState: vi.fn(async () => undefined), + }; + + (createFolder as Mock).mockResolvedValueOnce(mockFolder); + (checkFolderDuplicated as Mock).mockResolvedValueOnce({ + duplicatedFoldersResponse: [], + foldersWithDuplicates: [], + foldersWithoutDuplicates: [mockFolder], + }); + (uploadItemsParallelThunk as unknown as Mock).mockImplementation((thunkArgs) => thunkArgs); + const dispatchToUploadManager = vi.fn(({ files, parentFolderId, options, onFileUploadCallback }) => ({ + unwrap: () => + uploadFileWithManager({ + files: files.map((content: File, index: number) => ({ + taskId: `file-task-${index}`, + relatedTaskId: options.relatedTaskId, + filecontent: { content, name: content.name, size: content.size, type: 'txt', parentFolderId }, + userEmail: '', + parentFolderId, + })), + maxSpaceOccupiedCallback: vi.fn(), + uploadRepository: failingUploadRepository, + abortController: options.abortController, + options, + onFileUploadCallback, + }), + })); + const events: UploadFolderManagerEvents = { + onFolderUploadSuccess: vi.fn(), + onFolderUploadError: vi.fn(), + }; + + const folderUpload = uploadFoldersWithManager({ + payload: [ + { + currentFolderId: 'currentFolderId', + root: { + folderId: mockFolder.uuid, + childrenFiles: [new File(['a'], 'a.txt'), new File(['b'], 'b.txt')], + childrenFolders: [], + name: mockFolder.name, + fullPathEdited: 'path1', + }, + options: { taskId }, + }, + ], + selectedWorkspace: null, + dispatch: dispatchToUploadManager, + maxUploadFileSize: 100, + events, + }); + const hangTimeout = new Promise((_, reject) => + setTimeout(() => reject(new Error('folder upload never settled')), 2000), + ); + + await expect(Promise.race([folderUpload, hangTimeout])).resolves.toBeUndefined(); + expect(events.onFolderUploadSuccess).toHaveBeenCalledWith(taskId, { + folderName: 'MyFolder', + rootFolderUUID: mockFolder.uuid, + hasFailedFiles: true, + }); + expect(events.onFolderUploadError).not.toHaveBeenCalled(); + }); + test('When the folder itself cannot be created, then the failure is notified and no success is announced afterwards', async () => { const taskId = 'task-id'; diff --git a/src/app/network/UploadManager.test.ts b/src/app/network/UploadManager.test.ts index a38f202560..f32c699327 100644 --- a/src/app/network/UploadManager.test.ts +++ b/src/app/network/UploadManager.test.ts @@ -1,9 +1,9 @@ import { beforeEach, describe, expect, Mock, test, vi } from 'vitest'; -import { uploadFileWithManager, UploadManagerEvents } from './UploadManager'; +import { UploadFileParamsWithTaskId, uploadFileWithManager, UploadManagerEvents } from './UploadManager'; import errorService from 'services/error.service'; import { AppError } from '@internxt/sdk'; import uploadFile from 'app/drive/services/file.service/uploadFile'; -import DatabaseUploadRepository from 'app/repositories/DatabaseUploadRepository'; +import DatabaseUploadRepository, { PersistUploadRepository } from 'app/repositories/DatabaseUploadRepository'; import { DriveFileData } from 'app/drive/types'; import RetryManager, { RetryableTaskType } from './RetryManager'; import { TaskStatus } from 'app/tasks/types'; @@ -816,3 +816,155 @@ describe('uploadFileWithManager', () => { expect(openMaxSpaceOccupiedDialogMock).toHaveBeenCalledOnce(); }); }); + +describe('UploadManager settles every queued file', () => { + const HANG_TIMEOUT_MS = 2000; + const uploadStateError = new Error('IDB failure'); + + const failIfHangs = (promise: Promise): Promise => + Promise.race([ + promise, + new Promise((_, reject) => setTimeout(() => reject(new Error('upload never settled')), HANG_TIMEOUT_MS)), + ]); + + const buildFile = (name: string): UploadFileParamsWithTaskId => ({ + taskId: `${name}-task`, + filecontent: { + content: name as unknown as File, + type: 'text/plain', + name, + size: 1024, + parentFolderId: 'folder-1', + }, + userEmail: '', + parentFolderId: '', + }); + + const buildRepository = (failingUploadStateCalls: Record = {}): PersistUploadRepository => { + const callsByTaskId: Record = {}; + return { + setUploadState: vi.fn(async () => undefined), + removeUploadState: vi.fn(async () => undefined), + getUploadState: vi.fn(async (id: string) => { + callsByTaskId[id] = (callsByTaskId[id] ?? 0) + 1; + if (failingUploadStateCalls[id]?.includes(callsByTaskId[id])) throw uploadStateError; + return undefined; + }), + }; + }; + + const runUpload = ({ + files, + uploadRepository = buildRepository(), + events, + isUploadedFromFolder = false, + }: { + files: UploadFileParamsWithTaskId[]; + uploadRepository?: PersistUploadRepository; + events?: UploadManagerEvents; + isUploadedFromFolder?: boolean; + }) => + failIfHangs( + uploadFileWithManager({ + files, + maxSpaceOccupiedCallback: openMaxSpaceOccupiedDialogMock, + uploadRepository, + events, + options: { isUploadedFromFolder }, + }), + ); + + const uploadedNames = () => (uploadFile as Mock).mock.calls.map(([file]) => file.name); + + beforeEach(() => { + RetryManager.clearTasks(); + vi.clearAllMocks(); + vi.spyOn(errorService, 'castError').mockImplementation((e) => (e ?? new AppError('Unknown error')) as AppError); + vi.spyOn(errorService, 'reportError').mockReturnValue(); + (uploadFile as Mock).mockImplementation(async (file: { name: string }) => ({ ...mockFile1, name: file.name })); + }); + + test('when reading the upload state keeps failing, then the upload settles and the file is reported as failed', async () => { + const file = buildFile('file.txt'); + const events: UploadManagerEvents = { onUploadError: vi.fn() }; + const uploadRepository = buildRepository({ [file.taskId]: [1, 2] }); + + await expect(runUpload({ files: [file], uploadRepository, events })).rejects.toThrow(uploadStateError); + + expect(uploadFile).not.toHaveBeenCalled(); + expect(events.onUploadError).toHaveBeenCalledWith( + expect.objectContaining({ taskId: file.taskId }), + 'upload-failed', + ); + }); + + test('when reading the upload state fails on the retry after a 502, then the folder upload settles with the file failed and the others uploaded', async () => { + const failingFile = buildFile('failing.txt'); + const otherFiles = [buildFile('a.txt'), buildFile('b.txt')]; + (uploadFile as Mock).mockImplementationOnce(() => Promise.reject({ status: 502 })); + const uploadRepository = buildRepository({ [failingFile.taskId]: [2] }); + const events: UploadManagerEvents = { onUploadSuccess: vi.fn(), onUploadError: vi.fn() }; + + await expect( + runUpload({ files: [failingFile, ...otherFiles], uploadRepository, events, isUploadedFromFolder: true }), + ).rejects.toThrow(uploadStateError); + + expect(events.onUploadSuccess).toHaveBeenCalledTimes(otherFiles.length); + expect(events.onUploadError).toHaveBeenCalledWith( + expect.objectContaining({ taskId: failingFile.taskId }), + 'upload-failed', + ); + expect(RetryManager.getTasks()).toEqual([expect.objectContaining({ taskId: failingFile.taskId, retryable: true })]); + }); + + test('when the upload rejects without an error value, then the upload settles and the file is reported as failed', async () => { + (uploadFile as Mock).mockRejectedValue(undefined); + const events: UploadManagerEvents = { onUploadError: vi.fn() }; + + await expect(runUpload({ files: [buildFile('file.txt')], events })).rejects.toThrow('Unknown error'); + + expect(events.onUploadError).toHaveBeenCalledWith(expect.anything(), 'upload-failed'); + }); + + test('when a success listener throws, then every file is uploaded once and the queue keeps processing', async () => { + const files = ['a.txt', 'b.txt', 'c.txt'].map(buildFile); + const onFileUploadCallback = vi.fn(); + const events: UploadManagerEvents = { + onUploadSuccess: vi.fn((file) => { + if (file.taskId === files[0].taskId) throw new Error('listener failure'); + }), + }; + + const { uploadedFiles } = await failIfHangs( + uploadFileWithManager({ + files, + maxSpaceOccupiedCallback: openMaxSpaceOccupiedDialogMock, + uploadRepository: buildRepository(), + events, + onFileUploadCallback, + }), + ); + + expect(uploadedNames()).toEqual(['a.txt', 'b.txt', 'c.txt']); + expect(uploadedFiles.map(({ name }) => name)).toEqual(['a.txt', 'b.txt', 'c.txt']); + expect(onFileUploadCallback).toHaveBeenCalledTimes(files.length); + }); + + test('when a listener throws while the next queued file starts, then finished files are not uploaded again and the folder upload settles', async () => { + const files = Array.from({ length: 7 }, (_, i) => buildFile(`f${i}.txt`)); + const startingFile = files[6]; + const listenerError = new Error('listener failure'); + const events: UploadManagerEvents = { + onUploadStart: vi.fn((file) => { + if (file.taskId === startingFile.taskId) throw listenerError; + }), + onUploadSuccess: vi.fn(), + }; + + await expect(runUpload({ files, events, isUploadedFromFolder: true })).rejects.toThrow(listenerError); + + expect(uploadedNames()).toEqual(files.slice(0, 6).map((file) => file.filecontent.name)); + expect(events.onUploadSuccess).toHaveBeenCalledTimes(6); + expect(RetryManager.getTasks()).toEqual([expect.objectContaining({ taskId: startingFile.taskId })]); + }); +}); diff --git a/src/app/network/UploadManager.ts b/src/app/network/UploadManager.ts index 7e188f3818..5a2df33ae7 100644 --- a/src/app/network/UploadManager.ts +++ b/src/app/network/UploadManager.ts @@ -70,6 +70,8 @@ interface UploadFileWithManagerProps { events?: UploadManagerEvents; } +type UploadQueueCallback = (err: Error | null, res?: DriveFileData) => void; + export const uploadFileWithManager = (props: UploadFileWithManagerProps): Promise<{ uploadedFiles: DriveFileData[] }> => new UploadManager(props).run(); @@ -115,114 +117,14 @@ class UploadManager { }; private readonly uploadQueue: QueueObject = queue( - (fileData, next: (err: Error | null, res?: DriveFileData) => void) => { + (fileData, next: UploadQueueCallback) => { if (this.abortController?.signal.aborted ?? fileData.abortController?.signal.aborted) return; - this.manageMemoryUsage(); - - let uploadAttempts = 0; - const uploadId = randomBytes(10).toString('hex'); - const taskId = fileData.taskId; - this.uploadsProgress[uploadId] = 0; - - const file = fileData.filecontent; - - this.events?.onUploadStart?.(fileData, async () => { - if (this.abortController) this.abortController?.abort(); - else fileData?.abortController?.abort(); + const nextOnce = UploadManager.callOnce(next); + this.processFile(fileData, nextOnce).catch((error) => { + errorService.reportError(error); + nextOnce(errorService.castError(error)); }); - - const existsRelatedTask = !!fileData.relatedTaskId; - if (!existsRelatedTask) this.uploadRepository?.setUploadState(taskId, TaskStatus.InProcess); - const retryUploadType = fileData.fileType; - - const upload = async () => { - uploadAttempts++; - - this.events?.onUploadAttempt?.(fileData); - - const uploadStatus = this.uploadRepository?.getUploadState(fileData.relatedTaskId ?? taskId); - const isPaused = (await uploadStatus) === TaskStatus.Paused; - const continueUploadOptions = { - taskId: fileData.relatedTaskId ?? taskId, - isPaused, - isRetriedUpload: !!this.options?.isRetriedUpload, - }; - - let unsubscribeAbortListener: (() => void) | void; - - uploadFile( - { - name: file.name, - size: file.size, - type: retryUploadType ?? file.type, - content: file.content, - parentFolderId: file.parentFolderId, - }, - (uploadProgress) => { - this.uploadsProgress[uploadId] = uploadProgress; - this.events?.onUploadProgress?.(fileData, uploadProgress); - }, - { - isTeam: !!this.options?.ownerUserAuthenticationData?.workspaceId, - abortController: this.abortController ?? fileData.abortController, - ownerUserAuthenticationData: this.options?.ownerUserAuthenticationData, - abortCallback: (abort?: () => void) => { - unsubscribeAbortListener = this.events?.registerUploadAbort?.(fileData, () => abort?.()); - }, - isUploadedFromFolder: fileData.isUploadedFromFolder, - }, - continueUploadOptions, - ) - .then(async (driveFileData) => { - const isUploadAborted = this.abortController?.signal.aborted ?? fileData.abortController?.signal.aborted; - - if (isUploadAborted) { - throw new Error('Upload task cancelled'); - } - - const driveFileDataWithNameParsed = { ...driveFileData, name: file.name }; - this.filesUploadedList.push({ ...driveFileDataWithNameParsed, taskId }); - - await this.events?.onUploadSuccess?.(fileData, driveFileDataWithNameParsed); - - fileData.onFinishUploadFile?.(driveFileDataWithNameParsed, taskId); - - if (this.onFileUploadCallback) { - this.onFileUploadCallback(driveFileDataWithNameParsed); - } - next(null, driveFileDataWithNameParsed); - }) - .catch((error) => { - const isUploadAborted = - !!this.abortController?.signal.aborted || !!fileData.abortController?.signal.aborted || error === 'abort'; - const isLostConnectionError = - error instanceof ConnectionLostError || error.message === ErrorMessages.NetworkError; - const isNonRetryableError = UploadManager.isNonRetryableError(error); - - if ( - uploadAttempts < MAX_UPLOAD_ATTEMPTS && - !isUploadAborted && - !isLostConnectionError && - !isNonRetryableError - ) { - upload(); - } else { - this.handleUploadErrors({ - error, - fileData, - isUploadAborted, - isLostConnectionError, - next, - }); - } - }) - .finally(() => { - unsubscribeAbortListener?.(); - }); - }; - - upload(); }, this.filesGroups.small.concurrency, ); @@ -238,6 +140,143 @@ class UploadManager { this.events = props.events; } + private static callOnce(next: UploadQueueCallback): UploadQueueCallback { + let hasBeenCalled = false; + return (err, res) => { + if (hasBeenCalled) return; + hasBeenCalled = true; + next(err, res); + }; + } + + private async processFile(fileData: UploadFileParamsWithTaskId, next: UploadQueueCallback): Promise { + this.manageMemoryUsage(); + + const uploadId = randomBytes(10).toString('hex'); + this.uploadsProgress[uploadId] = 0; + + let driveFileData: DriveFileData; + try { + this.events?.onUploadStart?.(fileData, async () => { + if (this.abortController) this.abortController?.abort(); + else fileData?.abortController?.abort(); + }); + + const existsRelatedTask = !!fileData.relatedTaskId; + if (!existsRelatedTask) { + this.uploadRepository + ?.setUploadState(fileData.taskId, TaskStatus.InProcess) + .catch((error) => errorService.reportError(error)); + } + + driveFileData = await this.uploadWithRetries(fileData, uploadId); + } catch (error) { + this.handleUploadErrors({ error, fileData, next }); + return; + } + + const driveFileDataWithNameParsed = { ...driveFileData, name: fileData.filecontent.name }; + await this.notifyUploadSuccess(fileData, driveFileDataWithNameParsed); + next(null, driveFileDataWithNameParsed); + } + + private async uploadWithRetries( + fileData: UploadFileParamsWithTaskId, + uploadId: string, + uploadAttempts = 1, + ): Promise { + try { + return await this.uploadAttempt(fileData, uploadId); + } catch (error) { + const { isUploadAborted, isLostConnectionError } = this.classifyUploadError(error, fileData); + const shouldRetry = + uploadAttempts < MAX_UPLOAD_ATTEMPTS && + !isUploadAborted && + !isLostConnectionError && + !UploadManager.isNonRetryableError(error); + + if (!shouldRetry) throw error; + return this.uploadWithRetries(fileData, uploadId, uploadAttempts + 1); + } + } + + private async uploadAttempt(fileData: UploadFileParamsWithTaskId, uploadId: string): Promise { + const file = fileData.filecontent; + const taskId = fileData.relatedTaskId ?? fileData.taskId; + const abortListener: { unsubscribe?: (() => void) | void } = {}; + + try { + this.events?.onUploadAttempt?.(fileData); + + const isPaused = (await this.uploadRepository?.getUploadState(taskId)) === TaskStatus.Paused; + const continueUploadOptions = { + taskId, + isPaused, + isRetriedUpload: !!this.options?.isRetriedUpload, + }; + + const driveFileData = await uploadFile( + { + name: file.name, + size: file.size, + type: fileData.fileType ?? file.type, + content: file.content, + parentFolderId: file.parentFolderId, + }, + (uploadProgress) => { + this.uploadsProgress[uploadId] = uploadProgress; + this.events?.onUploadProgress?.(fileData, uploadProgress); + }, + { + isTeam: !!this.options?.ownerUserAuthenticationData?.workspaceId, + abortController: this.abortController ?? fileData.abortController, + ownerUserAuthenticationData: this.options?.ownerUserAuthenticationData, + abortCallback: (abort?: () => void) => { + abortListener.unsubscribe = this.events?.registerUploadAbort?.(fileData, () => abort?.()); + }, + isUploadedFromFolder: fileData.isUploadedFromFolder, + }, + continueUploadOptions, + ); + + const isUploadAborted = this.abortController?.signal.aborted ?? fileData.abortController?.signal.aborted; + if (isUploadAborted) { + throw new Error('Upload task cancelled'); + } + + return driveFileData; + } finally { + abortListener.unsubscribe?.(); + } + } + + private async notifyUploadSuccess(fileData: UploadFileParamsWithTaskId, driveFileData: DriveFileData): Promise { + this.filesUploadedList.push({ ...driveFileData, taskId: fileData.taskId }); + + await UploadManager.runListener(() => this.events?.onUploadSuccess?.(fileData, driveFileData)); + await UploadManager.runListener(() => fileData.onFinishUploadFile?.(driveFileData, fileData.taskId)); + await UploadManager.runListener(() => this.onFileUploadCallback?.(driveFileData)); + } + + private static async runListener(listener: () => unknown): Promise { + try { + await listener(); + } catch (error) { + errorService.reportError(error); + } + } + + private classifyUploadError( + error: unknown, + fileData: UploadFileParamsWithTaskId, + ): { isUploadAborted: boolean; isLostConnectionError: boolean } { + const isUploadAborted = + !!this.abortController?.signal.aborted || !!fileData.abortController?.signal.aborted || error === 'abort'; + const isLostConnectionError = + error instanceof ConnectionLostError || (error as Error | undefined)?.message === ErrorMessages.NetworkError; + return { isUploadAborted, isLostConnectionError }; + } + private static isNonRetryableError(error: unknown): boolean { return ( (error as { status?: number })?.status === HTTP_STATUS_CODES.PAYMENT_REQUIRED || @@ -278,16 +317,13 @@ class UploadManager { private handleUploadErrors({ error, fileData, - isUploadAborted, - isLostConnectionError, next, }: { error: unknown; fileData: UploadFileParamsWithTaskId; - isUploadAborted: boolean; - isLostConnectionError: boolean; - next: (err: Error | null, res?: DriveFileData) => void; + next: UploadQueueCallback; }) { + const { isUploadAborted, isLostConnectionError } = this.classifyUploadError(error, fileData); const castedError = errorService.castError(error); // Handle retry error if (castedError.message === 'Retryable file') { @@ -411,7 +447,7 @@ class UploadManager { failedUploadFiles.forEach((fileErrored) => { this.events?.onUploadError?.(fileErrored); - this.uploadRepository?.removeUploadState(fileErrored.taskId); + this.uploadRepository?.removeUploadState(fileErrored.taskId).catch((error) => errorService.reportError(error)); }); }