From 1783869b3c0f01b6e8a7fe659cff52c9119b81ae Mon Sep 17 00:00:00 2001 From: ShaniStrassProg <37214725715@mby.co.il> Date: Mon, 16 Sep 2024 12:04:07 +0300 Subject: [PATCH 1/2] save_in_memory --- .../DataAccess/ObjectManager.py | 35 +++++++++++-------- 1 file changed, 20 insertions(+), 15 deletions(-) diff --git a/Storage/NEW_KT_Storage/DataAccess/ObjectManager.py b/Storage/NEW_KT_Storage/DataAccess/ObjectManager.py index fea98ca2..80959b07 100644 --- a/Storage/NEW_KT_Storage/DataAccess/ObjectManager.py +++ b/Storage/NEW_KT_Storage/DataAccess/ObjectManager.py @@ -1,27 +1,32 @@ from typing import Dict, Any import json import sqlite3 -from KT_DB import ObjectManager +import os +import sys +sys.path.append(os.path.abspath(os.path.join(os.path.dirname(__file__), '..','..','..'))) +from DB.NEW_KT_DB.DataAccess.ObjectManager import ObjectManager as DBObjectManager from StorageManager import StorageManager +sys.path.append(os.path.abspath(os.path.join(os.path.dirname(__file__), '..'))) +from Models.MultipartUploadModel import MultipartUploadModel + class ObjectManager: - def __init__(self, db_file: str): + def __init__(self, db_file: str, storage_path=None): '''Initialize ObjectManager with the database connection.''' - self.object_manager = ObjectManager(db_file) - self.storage_manager = StorageManager(storage_path) - + self.object_manager = DBObjectManager(db_file) + self.storage_manager = StorageManager("storage_path") # for outer use: def save_in_memory(self, object): # insert object info into management table mng_{object_name}s # for exmple: object db_instance will be saved in table mng_db_instances - table_name = object_manager.convert_object_name_to_management_table_name(self.object_name) + table_name = str(object.__class__.__name__) - if not object_manager.is_management_table_exist(table_name): - object_manager.create_management_table(table_name) + if not self.object_manager.is_management_table_exist(table_name): + self.object_manager.create_management_table(table_name) - object_manager.insert_object_to_management_table(table_name, object) + self.object_manager.insert_object_to_management_table(table_name, object.to_sql()) def delete_from_memory(self,criteria='default'): @@ -30,9 +35,9 @@ def delete_from_memory(self,criteria='default'): if criteria == 'default': criteria = f'{self.pk_column} = {self.pk_value}' - table_name = object_manager.convert_object_name_to_management_table_name(self.object_name) + table_name = self.object_manager.convert_object_name_to_management_table_name(self.object_name) - object_manager.delete_data_from_table(table_name, criteria) + self.object_manager.delete_data_from_table(table_name, criteria) def update_in_memory(self, updates, criteria='default'): @@ -41,13 +46,13 @@ def update_in_memory(self, updates, criteria='default'): if criteria == 'default': criteria = f'{self.pk_column} = {self.pk_value}' - table_name = object_manager.convert_object_name_to_management_table_name(self.object_name) + table_name = self.object_manager.convert_object_name_to_management_table_name(self.object_name) - object_manager.update_object_in_management_table_by_criteria(table_name, updates, criteria) + self.object_manager.update_object_in_management_table_by_criteria(table_name, updates, criteria) def get_from_memory(self): - object_manager.get_object_from_management_table(self.object_id) + self.object_manager.get_object_from_management_table(self.object_id) def convert_object_attributes_to_dictionary(**kwargs): @@ -57,4 +62,4 @@ def convert_object_attributes_to_dictionary(**kwargs): for key, value in kwargs.items(): dict[key] = value - return dict + return dict \ No newline at end of file From db0485b38b23265348a2c9852507f190554acda6 Mon Sep 17 00:00:00 2001 From: ShaniStrassProg <37214725715@mby.co.il> Date: Mon, 16 Sep 2024 12:12:15 +0300 Subject: [PATCH 2/2] is_management_table_exist --- .../test_quota.cpython-312-pytest-8.3.2.pyc | Bin 7270 -> 7283 bytes DB/NEW_KT_DB/DataAccess/DBManager.py | 97 ++++--- DB/NEW_KT_DB/DataAccess/ObjectManager.py | 19 +- .../DataAccess/MultipartUploadManager.py | 114 +++++++++ .../Models/MultipartUploadModel.py | 42 +++ Storage/NEW_KT_Storage/Models/PartModel.py | 15 ++ .../Service/Classes/MultipartUploadService.py | 20 ++ .../tests/test_multipart_upload_manager.py | 242 ++++++++++++++++++ 8 files changed, 492 insertions(+), 57 deletions(-) create mode 100644 Storage/NEW_KT_Storage/DataAccess/MultipartUploadManager.py create mode 100644 Storage/NEW_KT_Storage/Models/MultipartUploadModel.py create mode 100644 Storage/NEW_KT_Storage/Models/PartModel.py create mode 100644 Storage/NEW_KT_Storage/Service/Classes/MultipartUploadService.py create mode 100644 Storage/NEW_KT_Storage/tests/test_multipart_upload_manager.py diff --git a/DB/DB_UserAdministration/tests/__pycache__/test_quota.cpython-312-pytest-8.3.2.pyc b/DB/DB_UserAdministration/tests/__pycache__/test_quota.cpython-312-pytest-8.3.2.pyc index 3c87af306a7b3468c36c446076dc0ffb82fd227d..58d1d836539c804a4259e7393aaabce332256a51 100644 GIT binary patch delta 78 zcmaE6@!5j=G%qg~0}#wO^>8D%8C delta 60 zcmext@yvqzG%qg~0}!lczrK;%jY;XNJCKuX6?1*&^~o{UCtja=ea7|W*XLbd65|~b P@0^ognzFfp>9ZsNyi^&X diff --git a/DB/NEW_KT_DB/DataAccess/DBManager.py b/DB/NEW_KT_DB/DataAccess/DBManager.py index dd0dde58..3886af23 100644 --- a/DB/NEW_KT_DB/DataAccess/DBManager.py +++ b/DB/NEW_KT_DB/DataAccess/DBManager.py @@ -8,25 +8,22 @@ def __init__(self, db_file: str): '''Initialize the database connection and create tables if they do not exist.''' self.connection = sqlite3.connect(db_file) - # rachel-8511, ShaniStrassProg def close(self): - '''Close the database connection.''' - self.connection.close() - + '''Close the database connection.''' + self.connection.close() # saraNoigershel def execute_query_with_multiple_results(self, query: str) -> Optional[List[Tuple]]: - '''Execute a given query and return the results.''' - try: - c = self.connection.cursor() - c.execute(query) - results = c.fetchall() - # self.connection.commit() ??? - return results if results else None - except OperationalError as e: - raise Exception(f'Error executing query {query}: {e}') - + '''Execute a given query and return the results.''' + try: + c = self.connection.cursor() + c.execute(query) + results = c.fetchall() + # self.connection.commit() ??? + return results if results else None + except OperationalError as e: + raise Exception(f'Error executing query {query}: {e}') # ShaniStrassProg def execute_query_with_single_result(self, query: str) -> Optional[Tuple]: @@ -40,7 +37,6 @@ def execute_query_with_single_result(self, query: str) -> Optional[Tuple]: except OperationalError as e: raise Exception(f'Error executing query {query}: {e}') - # Riki7649255 def execute_query_without_results(self, query: str): '''Execute a given query without waiting for any result.''' @@ -50,49 +46,54 @@ def execute_query_without_results(self, query: str): self.connection.commit() except OperationalError as e: raise Exception(f'Error executing query {query}: {e}') - # Yael, Riki7649255 def create_table(self, table_name, table_structure): '''create a table in a given db by given table_structure''' create_statement = f'''CREATE TABLE IF NOT EXISTS {table_name} ({table_structure})''' - execute_query_without_results(create_statement) - + self.execute_query_without_results(create_statement) # Riki7649255 based on rachel-8511, ShaniStrassProg def insert_data_into_table(self, table_name, data): insert_statement = f'''INSERT INTO {table_name} VALUES {data}''' - execute_query_without_results(insert_statement) + self.execute_query_without_results(insert_statement) + # # Riki7649255 based on rachel-8511, Shani + # def update_records_in_table(self, table_name: str, updates: Dict[str, Any], criteria: str) -> None: + # '''Update records in the specified table based on criteria.''' + + # # add documentation here + # set_clause = ', '.join([f'{k} = ?' for k in updates.keys()]) + # values = list(updates.values()) + + # update_statement = f''' + # UPDATE {table_name} + # SET {set_clause} + # WHERE {criteria} + # ''' + + # self.execute_query_without_results(update_statement) - # Riki7649255 based on rachel-8511, Shani def update_records_in_table(self, table_name: str, updates: Dict[str, Any], criteria: str) -> None: - '''Update records in the specified table based on criteria.''' - - # add documentation here - set_clause = ', '.join([f'{k} = ?' for k in updates.keys()]) - values = list(updates.values()) - - update_statement = f''' - UPDATE {table_name} - SET {set_clause} - WHERE {criteria} - ''' - - execute_query_without_results(update_statement) - + '''Update records in the specified table based on criteria.''' + set_clause = '(' + ', '.join([f'{k}' for k in updates.keys()]) + ') = (' + ', '.join([f'"{v}"' for v in updates.values()]) + ')' + update_statement = f''' + UPDATE {table_name} + SET {set_clause} + WHERE {criteria} + ''' + self.execute_query_without_results(update_statement) # Riki7649255 based on rachel-8511 def delete_data_from_table(self, table_name: str, criteria: str) -> None: - '''Delete a record from the specified table based on criteria.''' - - delete_statement = f''' - DELETE FROM {table_name} - WHERE {criteria} - ''') - - execute_query_without_results(delete_statement) - + '''Delete a record from the specified table based on criteria.''' + + delete_statement = f''' + DELETE FROM {table_name} + WHERE {criteria} + ''' + + self.execute_query_without_results(delete_statement) # rachel-8511, Riki7649255 def select_and_return_records_from_table(self, table_name: str, columns: List[str] = ['*'], criteria: str = '') -> Dict[int, Dict[str, Any]]: @@ -109,23 +110,21 @@ def select_and_return_records_from_table(self, table_name: str, columns: List[st if criteria: query += f' WHERE {criteria}' try: - results = execute_query_with_multiple_results(query) + results = self.execute_query_with_multiple_results(query) return {result[0]: dict(zip(columns, result[1:])) for result in results} except OperationalError as e: raise Exception(f'Error selecting from {table_name}: {e}') - # rachel-8511, ShaniStrassProg, Riki7649255 def describe_table(self, table_name: str) -> Dict[str, str]: '''Describe table structure.''' try: desc_statement = f'PRAGMA table_info({table_name})' - columns = execute_query_with_multiple_results(desc_statement) + columns = self.execute_query_with_multiple_results(desc_statement) return {col[1]: col[2] for col in columns} except OperationalError as e: raise Exception(f'Error describing table {table_name}: {e}') - - + # ShaniStrassProg # should be in ObjectManager and send the query to one of the execute_query functions # def is_json_column_contains_key_and_value(self, table_name: str, key: str, value: str) -> bool: @@ -144,7 +143,6 @@ def describe_table(self, table_name: str) -> Dict[str, str]: # print(f'Error: {e}') # return False - # Yael, ShaniStrassProg # should be in ObjectManager and send the query to one of the execute_query functions # def is_identifier_exist(self, table_name: str, value: str) -> bool: @@ -158,7 +156,6 @@ def describe_table(self, table_name: str) -> Dict[str, str]: # return c.fetchone()[0] > 0 # except sqlite3.OperationalError as e: # print(f'Error: {e}') - # sara-lea # should be in ObjectManager and send the query to one of the execute_query functions diff --git a/DB/NEW_KT_DB/DataAccess/ObjectManager.py b/DB/NEW_KT_DB/DataAccess/ObjectManager.py index 56c6948f..3f343eab 100644 --- a/DB/NEW_KT_DB/DataAccess/ObjectManager.py +++ b/DB/NEW_KT_DB/DataAccess/ObjectManager.py @@ -1,7 +1,10 @@ from typing import Dict, Any import json import sqlite3 -from DBManager import DBManager +import os +import sys +sys.path.append(os.path.abspath(os.path.join(os.path.dirname(__file__), '..'))) +from DataAccess.DBManager import DBManager class ObjectManager: def __init__(self, db_file: str): @@ -12,7 +15,7 @@ def __init__(self, db_file: str): # for internal use only: # Riki7649255 based on rachel-8511 - def create_management_table(self, table_name, table_structure='object_id INTEGER PRIMARY KEY AUTOINCREMENT,type_object TEXT NOT NULL,metadata TEXT NOT NULL') + def create_management_table(self, table_name, table_structure): self.db_manager.create_table(table_name, table_structure) @@ -42,13 +45,13 @@ def delete_object_from_management_table(self, table_name, criteria) -> None: self.db_manager.delete_data_from_table(table_name, criteria) - # rachel-8511, ShaniStrassProg is it needed? + # # rachel-8511, ShaniStrassProg is it needed? # def get_all_objects(self) -> Dict[int, Dict[str, Any]]: # '''Retrieve all objects from the database.''' # return self.db_manager.select(self.table_name, ['object_id', 'type_object', 'metadata']) - # rachel-8511 is it needed? + # # rachel-8511 is it needed? # def describe_table(self) -> Dict[str, str]: # '''Describe the schema of the table.''' # return self.db_manager.describe(self.table_name) @@ -58,9 +61,11 @@ def convert_object_name_to_management_table_name(object_name): return f'mng_{object_name}s' - def is_management_table_exist(table_name): - # check if table exists using single result query - return db_manager.execute_query_with_single_result(f'desc table {table_name}') + + def is_management_table_exist(self, table_name): + query = f"SELECT name FROM sqlite_master WHERE type='table' AND name='{table_name}'" + result = self.db_manager.execute_query_with_single_result(query) + return result is not None # for outer use: diff --git a/Storage/NEW_KT_Storage/DataAccess/MultipartUploadManager.py b/Storage/NEW_KT_Storage/DataAccess/MultipartUploadManager.py new file mode 100644 index 00000000..94404dc8 --- /dev/null +++ b/Storage/NEW_KT_Storage/DataAccess/MultipartUploadManager.py @@ -0,0 +1,114 @@ +import sqlite3 +import uuid +import os +import sys +import json +sys.path.append(os.path.abspath(os.path.join(os.path.dirname(__file__), '..'))) +from Models.MultipartUploadModel import MultipartUploadModel +""" +Imports the ObjectManager class from the KT_Cloud.DB.NEW_KT_DB.DataAccess module. + +The ObjectManager class is responsible for managing the storage and retrieval of objects in the database. +""" +from DataAccess.ObjectManager import ObjectManager + +class MultipartUploadManager: + def __init__(self, db_file: str, storage_path: str = None): + self.storage_path = storage_path + self.object_manager = ObjectManager(db_file) + self.db_manager = self.object_manager.object_manager + self.create_table() + + def create_table(self): + table_schema = 'upload_id TEXT PRIMARY KEY, bucket_name TEXT NOT NULL, object_name TEXT NOT NULL, parts TEXT NOT NULL' + self.db_manager.create_management_table('MultipartUploadModel', table_schema) + + + + def create_multipart_upload(self, multipart_upload: MultipartUploadModel) -> str: + """יוצר תהליך העלאת חלקים ומחזיר UploadId ייחודי""" + if not isinstance(multipart_upload, MultipartUploadModel): + raise TypeError('Expected an instance of MultipartUploadModel') + self.object_manager.save_in_memory(multipart_upload) + return multipart_upload.upload_id + + # def create_multipart_upload(self, multipart_upload: MultipartUploadModel) -> str: + # """יוצר תהליך העלאת חלקים ומחזיר UploadId ייחודי""" + # # multipart_upload.upload_id = str(uuid.uuid4()) + # multipart_dict = multipart_upload.to_dict() + # multipart_dict['parts'] = json.dumps(multipart_dict['parts']) + # self.db_manager.insert_object_to_management_table(self.table_name, multipart_dict) + # return multipart_upload.upload_id + # def upload_part(self, bucket_name: str, object_key: str, upload_id: str, new_part: PartModel) -> str: + # """מעלה חלק מסוים עבור bucket ואובייקט""" + # part_file_path = os.path.join(self.storage_path, f'{bucket_name}/{object_key}_part_{new_part.part_number}') + # self.object_manager.create_file(part_file_path, new_part.body) + + # # עדכון אובייקט ההעלאה עם החלק החדש + # obj_parts = self.db_manager.get_object_from_management_table(upload_id) + # multipart_upload = MultipartUploadModel(obj_parts['bucket_name'], obj_parts['object_key']) + # multipart_upload.parts = obj_parts['parts'] + # multipart_upload.parts.append({ + # 'PartNumber': new_part.part_number, + # 'FilePath': part_file_path + # }) + + # # עדכון במסד הנתונים + # self.db_manager.update_object_in_management_table_by_criteria(self.table_name, multipart_upload, upload_id) + + # # Generate a fake ETag for the example + # return f'etag_{new_part.part_number}' + + # def complete_multipart_upload(self, bucket_name: str, object_key: str, upload_id: str): + # """משלים את ההעלאה, מאחד את כל החלקים לקובץ אחד""" + # obj_parts = self.db_manager.get_object_from_management_table(upload_id) + # complete_file_path = os.path.join(self.storage_path, f'{bucket_name}/{object_key}_complete') + + # with open(complete_file_path, 'wb') as complete_file: + # for part in sorted(obj_parts.parts, key=lambda x: x['PartNumber']): + # part_file_path = part['FilePath'] + # with open(part_file_path, 'rb') as part_file: + # complete_file.write(part_file.read()) + + # # setObject(bucket_name, key, complete_file_path) + + # # מחיקת חלקי הקבצים + # for part in obj_parts.parts: + # os.remove(part['FilePath']) + import sqlite3 + +def select_all_from_table(db_file: str, table_name: str): + """Selects all data from a specified table in the database.""" + try: + # יצירת חיבור למסד הנתונים + conn = sqlite3.connect(db_file) + cursor = conn.cursor() + + # הרצת שאילתת SELECT + query = f"SELECT * FROM {table_name}" + cursor.execute(query) + + # שליפת כל התוצאות + rows = cursor.fetchall() + + # סגירת החיבור למסד הנתונים + cursor.close() + conn.close() + + return rows + except sqlite3.Error as e: + print(f"An error occurred: {e}") + return None + +MultipartUploadModel_new = MultipartUploadModel("my_bucket", "my_object") +print(MultipartUploadModel_new.__class__.__name__) + +MultipartUploadManager_new = MultipartUploadManager("my_db.db", "my_storage_path") +MultipartUploadModel_new.upload_id = MultipartUploadManager_new.create_multipart_upload(MultipartUploadModel_new) +db_file = "my_db.db" +table_name = "MultipartUploadModel" +# select_and_return_records_from_table# קריאה לפונקציה +all_data = select_all_from_table(db_file, table_name) + +# הדפסת התוצאות +print(all_data) \ No newline at end of file diff --git a/Storage/NEW_KT_Storage/Models/MultipartUploadModel.py b/Storage/NEW_KT_Storage/Models/MultipartUploadModel.py new file mode 100644 index 00000000..7fe5839c --- /dev/null +++ b/Storage/NEW_KT_Storage/Models/MultipartUploadModel.py @@ -0,0 +1,42 @@ +import uuid +from typing import Optional +from datetime import datetime +import json + + +class MultipartUploadModel: + def __init__(self, bucket_name: str, object_key: str): + """Create a unique upload model with bucket and object key""" + self.upload_id = str(uuid.uuid4()) # Create a unique identifier for the upload + self.bucket_name = bucket_name # Bucket name + self.object_key = object_key # Object key (usually file name) + self.parts = "[] " # List of parts + + def to_sql(self): + # Convert the model instance to a dictionary + data_dict = self.to_dict() + values = '(' + ", ".join(f'\'{json.dumps(v)}\'' if isinstance(v, dict) or isinstance(v, list) else f'\'{v}\'' if isinstance(v, str) else f'\'{str(v)}\'' + for v in data_dict.values()) + ')' + return values + + def to_dict(self): + return { + 'bucket_name': self.bucket_name, + 'object_key': self.object_key, + 'upload_id': self.upload_id, + 'parts': self.parts + } + # def __init__(self, bucket_name: str, object_key: str): + # """יצירת מודל העלאה עם bucket ומפתח אובייקט ייחודי""" + # self.bucket_name = bucket_name # שם ה-bucket + # self.object_key = object_key # מפתח האובייקט (לרוב שם קובץ) + # self.upload_id = str(uuid.uuid4()) # יצירת מזהה ייחודי עבור ההעלאה + # self.parts = [] # רשימת חלקים + + # def to_dict(self): + # return { + # 'bucket_name': self.bucket_name, + # 'object_key': self.object_key, + # 'upload_id': self.upload_id, + # 'parts': self.parts + # } \ No newline at end of file diff --git a/Storage/NEW_KT_Storage/Models/PartModel.py b/Storage/NEW_KT_Storage/Models/PartModel.py new file mode 100644 index 00000000..9eb20036 --- /dev/null +++ b/Storage/NEW_KT_Storage/Models/PartModel.py @@ -0,0 +1,15 @@ +class PartModel: + def __init__(self, part_number: int, body: str, etag: str = None, last_modified: Optional[datetime] = None): + """יצירת מודל של חלק מסוים עם מספר, גוף תוכן ו-ETag""" + self.part_number = part_number + self.etag = etag + self.last_modified = last_modified + # self.body = body + + def to_dict(self) -> dict: + return { + 'PartNumber': self.part_number, + 'ETag': self.etag, + 'LastModified': self.last_modified.isoformat() if self.last_modified else None, + # 'body': self.body + } \ No newline at end of file diff --git a/Storage/NEW_KT_Storage/Service/Classes/MultipartUploadService.py b/Storage/NEW_KT_Storage/Service/Classes/MultipartUploadService.py new file mode 100644 index 00000000..5e9fadff --- /dev/null +++ b/Storage/NEW_KT_Storage/Service/Classes/MultipartUploadService.py @@ -0,0 +1,20 @@ +from Models import MultipartUploadModel , PartModel +from DataAccess import MultipartUploadManager + +class MultipartUploadService: + def __init__(self): + self.manager = MultipartUploadManager() + + def create_multipart_upload(self, bucket_name: str, object_key: str) -> str: + """יוזם את תהליך ההעלאה עבור bucket ואובייקט""" + multipart_upload = MultipartUploadModel(bucket_name=bucket_name, object_key=object_key) + return self.manager.create_multipart_upload(multipart_upload) + + def upload_part(self, bucket_name: str, object_key: str, upload_id: str, part_number: int, body: bytes) -> str: + """מעלה חלק של קובץ עבור bucket ואובייקט ומחזיר את ETag של החלק""" + new_part = PartModel(part_number=part_number, body=body) + return self.manager.upload_part(bucket_name, object_key, upload_id, new_part) + + def complete_upload(self, bucket_name: str, object_key: str, upload_id: str): + """משלים את תהליך ההעלאה""" + self.manager.complete_multipart_upload(bucket_name, object_key, upload_id) \ No newline at end of file diff --git a/Storage/NEW_KT_Storage/tests/test_multipart_upload_manager.py b/Storage/NEW_KT_Storage/tests/test_multipart_upload_manager.py new file mode 100644 index 00000000..6529ebc7 --- /dev/null +++ b/Storage/NEW_KT_Storage/tests/test_multipart_upload_manager.py @@ -0,0 +1,242 @@ +import os +import pytest +from unittest.mock import Mock, patch +from DataAccess.MultipartUploadManager import MultipartUploadManager + +@pytest.fixture +def multipart_upload_manager(): + db_manager_mock = Mock() + storage_path = '/tmp/test_storage' + return MultipartUploadManager(db_manager_mock, storage_path) + +def test_complete_multipart_upload_success(multipart_upload_manager, tmp_path): + bucket_name = 'test-bucket' + object_key = 'test-object' + upload_id = 'test-upload-id' + + # Create test part files + part1 = tmp_path / 'part1' + part2 = tmp_path / 'part2' + part1.write_bytes(b'Part 1 content') + part2.write_bytes(b'Part 2 content') + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = { + 'parts': [ + {'PartNumber': 1, 'FilePath': str(part1)}, + {'PartNumber': 2, 'FilePath': str(part2)} + ] + } + + multipart_upload_manager.storage_path = str(tmp_path) + + multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + complete_file_path = os.path.join(str(tmp_path), f'{bucket_name}/{object_key}_complete') + assert os.path.exists(complete_file_path) + with open(complete_file_path, 'rb') as f: + assert f.read() == b'Part 1 contentPart 2 content' + + assert not os.path.exists(str(part1)) + assert not os.path.exists(str(part2)) + +def test_complete_multipart_upload_empty_parts(multipart_upload_manager, tmp_path): + bucket_name = 'test-bucket' + object_key = 'test-object' + upload_id = 'test-upload-id' + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = { + 'parts': [] + } + + multipart_upload_manager.storage_path = str(tmp_path) + + multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + complete_file_path = os.path.join(str(tmp_path), f'{bucket_name}/{object_key}_complete') + assert os.path.exists(complete_file_path) + assert os.path.getsize(complete_file_path) == 0 + +def test_complete_multipart_upload_file_not_found(multipart_upload_manager, tmp_path): + bucket_name = 'test-bucket' + object_key = 'test-object' + upload_id = 'test-upload-id' + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = { + 'parts': [ + {'PartNumber': 1, 'FilePath': '/non/existent/path'} + ] + } + + multipart_upload_manager.storage_path = str(tmp_path) + + with pytest.raises(FileNotFoundError): + multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + +def test_complete_multipart_upload_unsorted_parts(multipart_upload_manager, tmp_path): + bucket_name = 'test-bucket' + object_key = 'test-object' + upload_id = 'test-upload-id' + + part1 = tmp_path / 'part1' + part2 = tmp_path / 'part2' + part1.write_bytes(b'Part 1 content') + part2.write_bytes(b'Part 2 content') + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = { + 'parts': [ + {'PartNumber': 2, 'FilePath': str(part2)}, + {'PartNumber': 1, 'FilePath': str(part1)} + ] + } + + multipart_upload_manager.storage_path = str(tmp_path) + + multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + complete_file_path = os.path.join(str(tmp_path), f'{bucket_name}/{object_key}_complete') + assert os.path.exists(complete_file_path) + with open(complete_file_path, 'rb') as f: + assert f.read() == b'Part 1 contentPart 2 content' + +def test_complete_multipart_upload_permission_error(multipart_upload_manager, tmp_path): + bucket_name = 'test-bucket' + object_key = 'test-object' + upload_id = 'test-upload-id' + + part1 = tmp_path / 'part1' + part1.write_bytes(b'Part 1 content') + os.chmod(str(part1), 0o000) + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = { + 'parts': [ + {'PartNumber': 1, 'FilePath': str(part1)} + ] + } + + multipart_upload_manager.storage_path = str(tmp_path) + + with pytest.raises(PermissionError): + multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + os.chmod(str(part1), 0o644) # Restore permissions for cleanup +class TestMultipartUploadManager: + @pytest.fixture + def multipart_upload_manager(self, tmp_path): + storage_path = tmp_path / "storage" + storage_path.mkdir() + db_manager = MagicMock() + return MultipartUploadManager(str(storage_path), db_manager) + + @pytest.mark.asyncio + async def test_complete_multipart_upload_success(self, multipart_upload_manager, tmp_path): + bucket_name = "test-bucket" + object_key = "test-object" + upload_id = "test-upload-id" + + # Create test part files + part1 = tmp_path / "part1" + part2 = tmp_path / "part2" + part1.write_bytes(b"Part 1 content") + part2.write_bytes(b"Part 2 content") + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = MagicMock( + parts=[ + {"PartNumber": 1, "FilePath": str(part1)}, + {"PartNumber": 2, "FilePath": str(part2)}, + ] + ) + + await multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + complete_file_path = os.path.join(multipart_upload_manager.storage_path, f'{bucket_name}/{object_key}_complete') + assert os.path.exists(complete_file_path) + with open(complete_file_path, 'rb') as f: + assert f.read() == b"Part 1 contentPart 2 content" + + assert not os.path.exists(str(part1)) + assert not os.path.exists(str(part2)) + + @pytest.mark.asyncio + async def test_complete_multipart_upload_empty_parts(self, multipart_upload_manager): + bucket_name = "test-bucket" + object_key = "test-object" + upload_id = "test-upload-id" + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = MagicMock(parts=[]) + + with pytest.raises(ValueError, match="No parts found for the multipart upload"): + await multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + @pytest.mark.asyncio + async def test_complete_multipart_upload_missing_part_file(self, multipart_upload_manager, tmp_path): + bucket_name = "test-bucket" + object_key = "test-object" + upload_id = "test-upload-id" + + part1 = tmp_path / "part1" + part1.write_bytes(b"Part 1 content") + non_existent_part = tmp_path / "non_existent_part" + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = MagicMock( + parts=[ + {"PartNumber": 1, "FilePath": str(part1)}, + {"PartNumber": 2, "FilePath": str(non_existent_part)}, + ] + ) + + with pytest.raises(FileNotFoundError): + await multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + @pytest.mark.asyncio + async def test_complete_multipart_upload_out_of_order_parts(self, multipart_upload_manager, tmp_path): + bucket_name = "test-bucket" + object_key = "test-object" + upload_id = "test-upload-id" + + part1 = tmp_path / "part1" + part2 = tmp_path / "part2" + part1.write_bytes(b"Part 1 content") + part2.write_bytes(b"Part 2 content") + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = MagicMock( + parts=[ + {"PartNumber": 2, "FilePath": str(part2)}, + {"PartNumber": 1, "FilePath": str(part1)}, + ] + ) + + await multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + complete_file_path = os.path.join(multipart_upload_manager.storage_path, f'{bucket_name}/{object_key}_complete') + assert os.path.exists(complete_file_path) + with open(complete_file_path, 'rb') as f: + assert f.read() == b"Part 1 contentPart 2 content" + + @pytest.mark.asyncio + async def test_complete_multipart_upload_large_number_of_parts(self, multipart_upload_manager, tmp_path): + bucket_name = "test-bucket" + object_key = "test-object" + upload_id = "test-upload-id" + + num_parts = 1000 + parts = [] + expected_content = b"" + + for i in range(1, num_parts + 1): + part_file = tmp_path / f"part{i}" + part_content = f"Part {i} content".encode() + part_file.write_bytes(part_content) + parts.append({"PartNumber": i, "FilePath": str(part_file)}) + expected_content += part_content + + multipart_upload_manager.db_manager.get_object_from_management_table.return_value = MagicMock(parts=parts) + + await multipart_upload_manager.complete_multipart_upload(bucket_name, object_key, upload_id) + + complete_file_path = os.path.join(multipart_upload_manager.storage_path, f'{bucket_name}/{object_key}_complete') + assert os.path.exists(complete_file_path) + with open(complete_file_path, 'rb') as f: + assert f.read() == expected_content + + for part in parts: + assert not os.path.exists(part['FilePath'])