From 448ba239d7e790d56a141063948d7ca6986100d2 Mon Sep 17 00:00:00 2001 From: Hugo Antunes Date: Fri, 24 Nov 2023 12:25:12 +0100 Subject: [PATCH 1/5] PILOT-4152: Update cli logic to fit dynamic chunk size when uploading --- app/services/file_manager/file_upload/upload_client.py | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/app/services/file_manager/file_upload/upload_client.py b/app/services/file_manager/file_upload/upload_client.py index bea34f29..09a3d37b 100644 --- a/app/services/file_manager/file_upload/upload_client.py +++ b/app/services/file_manager/file_upload/upload_client.py @@ -61,7 +61,7 @@ def __init__( self.user = UserConfig() self.operator = self.user.username self.upload_message = upload_message - self.chunk_size = AppConfig.Env.chunk_size # remove + self.chunk_size = AppConfig.Env.chunk_size self.base_url = { AppConfig.Env.green_zone: AppConfig.Connections.url_upload_greenroom, AppConfig.Env.core_zone: AppConfig.Connections.url_upload_core, @@ -111,7 +111,7 @@ def resume_upload(self, unfinished_file_objects: List[FileObject]) -> List[FileO Parameter: - unfinished_file_objects(List[FileObject]): the unfinished items that need to be resumed. return: - - list of FileObject: the infomation retrieved from backend. + - list of FileObject: the information retrieved from backend. - resumable_id(str): the unique identifier for multipart upload. - object_path(str): the path in the object storage. - local_path(str): the local path of file. @@ -121,6 +121,7 @@ def resume_upload(self, unfinished_file_objects: List[FileObject]) -> List[FileO headers = {'Authorization': 'Bearer ' + self.user.access_token, 'Session-ID': self.user.session_id} url = AppConfig.Connections.url_bff + f'/v1/project/{self.project_code}/files/resumable' rid_file_object_map = {x.resumable_id: x for x in unfinished_file_objects} + payload = { 'bucket': self.bucket, 'zone': self.zone, @@ -129,6 +130,9 @@ def resume_upload(self, unfinished_file_objects: List[FileObject]) -> List[FileO 'object_path': file_object.object_path, 'item_id': file_object.item_id, 'resumable_id': file_object.resumable_id, + 'chunks_info': { + '': {'etag': file_object.get('ETag'), 'chunk_size': file_object.get('ChunkSize')} + }, } for file_object in unfinished_file_objects ], @@ -353,6 +357,7 @@ def upload_chunk(self, file_object: FileObject, chunk_number: int, chunk: str, e 'key': file_object.item_id, 'upload_id': file_object.resumable_id, 'chunk_number': chunk_number, + 'chunk_size': self.chunk_size, } headers = { 'Authorization': 'Bearer ' + self.user.access_token, From bb374ea64555e219c6c9d43a70165bcf569615a0 Mon Sep 17 00:00:00 2001 From: Hugo Antunes Date: Thu, 30 Nov 2023 13:54:13 +0100 Subject: [PATCH 2/5] PILOT-4152: set chunk_size to upload_chunk in stream_upload --- .../file_manager/file_upload/upload_client.py | 19 +++++++++++-------- 1 file changed, 11 insertions(+), 8 deletions(-) diff --git a/app/services/file_manager/file_upload/upload_client.py b/app/services/file_manager/file_upload/upload_client.py index 09a3d37b..138f3ed6 100644 --- a/app/services/file_manager/file_upload/upload_client.py +++ b/app/services/file_manager/file_upload/upload_client.py @@ -130,9 +130,6 @@ def resume_upload(self, unfinished_file_objects: List[FileObject]) -> List[FileO 'object_path': file_object.object_path, 'item_id': file_object.item_id, 'resumable_id': file_object.resumable_id, - 'chunks_info': { - '': {'etag': file_object.get('ETag'), 'chunk_size': file_object.get('ChunkSize')} - }, } for file_object in unfinished_file_objects ], @@ -304,7 +301,13 @@ def stream_upload(self, file_object: FileObject, pool: ThreadPool) -> List[Apply # after all the chunks have been uploaded. chunk_result = [] while True: - chunk = f.read(self.chunk_size) + + if file_object.uploaded_chunks: + chunk_size = file_object.uploaded_chunks[0].get('chunk_size') + else: + chunk_size = self.chunk_size + + chunk = f.read(chunk_size) chunk_etag = file_object.uploaded_chunks.get(str(count + 1)) local_chunk_etag = hashlib.md5(chunk).hexdigest() if not chunk: @@ -316,11 +319,11 @@ def stream_upload(self, file_object: FileObject, pool: ThreadPool) -> List[Apply if chunk_etag != local_chunk_etag: SrvErrorHandler.customized_handle(ECustomizedError.INVALID_CHUNK_UPLOAD, value=count + 1) raise INVALID_CHUNK_ETAG(count + 1) - file_object.update_progress(self.chunk_size) + file_object.update_progress(chunk_size) else: res = pool.apply_async( self.upload_chunk, - args=(file_object, count + 1, chunk, local_chunk_etag), + args=(file_object, count + 1, chunk, local_chunk_etag, chunk_size), ) chunk_result.append(res) @@ -330,7 +333,7 @@ def stream_upload(self, file_object: FileObject, pool: ThreadPool) -> List[Apply return chunk_result - def upload_chunk(self, file_object: FileObject, chunk_number: int, chunk: str, etag: str) -> None: + def upload_chunk(self, file_object: FileObject, chunk_number: int, chunk: str, etag: str, chunk_size: int) -> None: """ Summary: The function is to upload a chunk directly into minio storage. @@ -357,7 +360,7 @@ def upload_chunk(self, file_object: FileObject, chunk_number: int, chunk: str, e 'key': file_object.item_id, 'upload_id': file_object.resumable_id, 'chunk_number': chunk_number, - 'chunk_size': self.chunk_size, + 'chunk_size': chunk_size, } headers = { 'Authorization': 'Bearer ' + self.user.access_token, From e7605014d843b39d0372284da98ccf57e9bab038 Mon Sep 17 00:00:00 2001 From: Hugo Antunes Date: Thu, 30 Nov 2023 13:56:17 +0100 Subject: [PATCH 3/5] PILOT-4152: fix tests --- .../app/services/file_manager/file_upload/test_upload_client.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/app/services/file_manager/file_upload/test_upload_client.py b/tests/app/services/file_manager/file_upload/test_upload_client.py index d1809d0b..53c3cd76 100644 --- a/tests/app/services/file_manager/file_upload/test_upload_client.py +++ b/tests/app/services/file_manager/file_upload/test_upload_client.py @@ -73,7 +73,7 @@ def test_chunk_upload(httpx_mock, mocker): mocker.patch('app.services.file_manager.file_upload.models.FileObject.generate_meta', return_value=(1, 1)) test_obj = FileObject('test', 'test', 'test', 'test', 'test') - res = upload_client.upload_chunk(test_obj, 0, b'1', 'test_etag') + res = upload_client.upload_chunk(test_obj, 0, b'1', 'test_etag', 10) assert test_obj.progress_bar.n == 1 assert res.status_code == 200 From ec9a6ef5ec4e76f6f61e495bc49d3fdba521199e Mon Sep 17 00:00:00 2001 From: Hugo Antunes Date: Thu, 30 Nov 2023 15:52:53 +0100 Subject: [PATCH 4/5] PILOT-4152: get chunk_size by chunk --- app/services/file_manager/file_upload/models.py | 5 +++-- app/services/file_manager/file_upload/upload_client.py | 9 +++------ 2 files changed, 6 insertions(+), 8 deletions(-) diff --git a/app/services/file_manager/file_upload/models.py b/app/services/file_manager/file_upload/models.py index 8d0b481b..aaf36ffb 100644 --- a/app/services/file_manager/file_upload/models.py +++ b/app/services/file_manager/file_upload/models.py @@ -7,7 +7,8 @@ from os.path import basename from os.path import dirname from os.path import getsize -from typing import List +from typing import Any +from typing import Dict from typing import Tuple from tqdm import tqdm @@ -62,7 +63,7 @@ class FileObject: total_chunks: int # resumable info - uploaded_chunks: List[dict] + uploaded_chunks: Dict[str, Dict[str, Any]] # progress bar object progress_bar = None diff --git a/app/services/file_manager/file_upload/upload_client.py b/app/services/file_manager/file_upload/upload_client.py index 138f3ed6..6c28d7e3 100644 --- a/app/services/file_manager/file_upload/upload_client.py +++ b/app/services/file_manager/file_upload/upload_client.py @@ -301,14 +301,11 @@ def stream_upload(self, file_object: FileObject, pool: ThreadPool) -> List[Apply # after all the chunks have been uploaded. chunk_result = [] while True: - - if file_object.uploaded_chunks: - chunk_size = file_object.uploaded_chunks[0].get('chunk_size') - else: - chunk_size = self.chunk_size + chunk = file_object.uploaded_chunks.get(str(count + 1)) + chunk_etag = chunk.get('etag') + chunk_size = chunk.get('chunk_size', self.chunk_size) chunk = f.read(chunk_size) - chunk_etag = file_object.uploaded_chunks.get(str(count + 1)) local_chunk_etag = hashlib.md5(chunk).hexdigest() if not chunk: break From ceb5e475c78462bd5173e3c4495e227f6ceab447 Mon Sep 17 00:00:00 2001 From: Hugo Antunes Date: Thu, 30 Nov 2023 16:52:45 +0100 Subject: [PATCH 5/5] add default value --- app/services/file_manager/file_upload/upload_client.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/services/file_manager/file_upload/upload_client.py b/app/services/file_manager/file_upload/upload_client.py index 6c28d7e3..0640f3e5 100644 --- a/app/services/file_manager/file_upload/upload_client.py +++ b/app/services/file_manager/file_upload/upload_client.py @@ -301,7 +301,7 @@ def stream_upload(self, file_object: FileObject, pool: ThreadPool) -> List[Apply # after all the chunks have been uploaded. chunk_result = [] while True: - chunk = file_object.uploaded_chunks.get(str(count + 1)) + chunk = file_object.uploaded_chunks.get(str(count + 1), {}) chunk_etag = chunk.get('etag') chunk_size = chunk.get('chunk_size', self.chunk_size)