From 189e6d80b0a3b9f3ba11b7211113066a55614f51 Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Thu, 18 Dec 2025 18:52:27 +0530 Subject: [PATCH 1/9] feat:StartUploadSessionMethod --- docs/api_reference/file.rst | 1 + docs/getting_started.rst | 3 ++- nisystemlink/clients/file/_file_client.py | 16 ++++++++++++++ nisystemlink/clients/file/models/__init__.py | 1 + .../models/_upload_session_start_response.py | 21 +++++++++++++++++++ tests/integration/file/test_file_client.py | 17 +++++++++++++++ 6 files changed, 58 insertions(+), 1 deletion(-) create mode 100644 nisystemlink/clients/file/models/_upload_session_start_response.py diff --git a/docs/api_reference/file.rst b/docs/api_reference/file.rst index 058c5002..4ef05fc2 100644 --- a/docs/api_reference/file.rst +++ b/docs/api_reference/file.rst @@ -15,6 +15,7 @@ nisystemlink.clients.file .. automethod:: upload_file .. automethod:: download_file .. automethod:: update_metadata + .. automethod:: start_upload_session .. automodule:: nisystemlink.clients.file.models :members: diff --git a/docs/getting_started.rst b/docs/getting_started.rst index 4608f959..bc6fc15e 100644 --- a/docs/getting_started.rst +++ b/docs/getting_started.rst @@ -203,7 +203,8 @@ default connection. The default connection depends on your environment. With a :class:`.FileClient` object, you can: -* Get the list of files, download and delete files +* Get the list of files, download and delete files. +* Start upload sessions for chunked file uploads. Examples ~~~~~~~~ diff --git a/nisystemlink/clients/file/_file_client.py b/nisystemlink/clients/file/_file_client.py index 052091fd..c26a8a21 100644 --- a/nisystemlink/clients/file/_file_client.py +++ b/nisystemlink/clients/file/_file_client.py @@ -275,3 +275,19 @@ def update_metadata(self, metadata: models.UpdateMetadataRequest, id: str) -> No Raises: ApiException: if unable to communicate with the File Service. """ + + @post("service-groups/Default/upload-sessions/start", args=[Query(name="workspace")]) + def start_upload_session( + self, workspace: str | None = None + ) -> models.UploadSessionStartResponse: + """Start an upload session for uploading a file in chunks. + + Args: + workspace: The id of the workspace the file belongs to. Defaults to None. + + Returns: + Upload session information including the session ID. + + Raises: + ApiException: if unable to communicate with the File Service. + """ diff --git a/nisystemlink/clients/file/models/__init__.py b/nisystemlink/clients/file/models/__init__.py index 4d1dede8..2580cfa6 100644 --- a/nisystemlink/clients/file/models/__init__.py +++ b/nisystemlink/clients/file/models/__init__.py @@ -5,5 +5,6 @@ from ._operations import V1Operations from ._update_metadata import UpdateMetadataRequest from ._file_linq_query import FileLinqQueryRequest, FileLinqQueryResponse +from ._upload_session_start_response import UploadSessionStartResponse # flake8: noqa diff --git a/nisystemlink/clients/file/models/_upload_session_start_response.py b/nisystemlink/clients/file/models/_upload_session_start_response.py new file mode 100644 index 00000000..3d3396cc --- /dev/null +++ b/nisystemlink/clients/file/models/_upload_session_start_response.py @@ -0,0 +1,21 @@ +from datetime import datetime + +from nisystemlink.clients.core._uplink._json_model import JsonModel + + +class UploadSessionStartResponse(JsonModel): + """Response model for starting an upload session.""" + + id: str + """ + The session id created. + + example: 54837669-8cf5-469e-bf7d-26cb808c8f24 + """ + + created_at: datetime + """ + The date and time the upload session has started. + + example: 2018-05-15T18:54:27.519Z + """ diff --git a/tests/integration/file/test_file_client.py b/tests/integration/file/test_file_client.py index 4ae5d7e5..810374a2 100644 --- a/tests/integration/file/test_file_client.py +++ b/tests/integration/file/test_file_client.py @@ -263,3 +263,20 @@ def test__query_files_linq__filter_returns_no_results(self, client: FileClient): assert response.total_count is not None assert response.total_count.value == 0 assert response.total_count.relation == "eq" + + def test__start_upload_session__returns_session_id(self, client: FileClient): + response = client.start_upload_session() + + assert response is not None + assert response.id is not None + assert isinstance(response.id, str) + assert response.created_at is not None + assert isinstance(response.created_at, datetime) + + def test__start_upload_session__with_invalid_workspace__raises( + self, client: FileClient + ): + invalid_workspace_id = "invalid-workspace-id" + + with pytest.raises(ApiException): + client.start_upload_session(workspace=invalid_workspace_id) From 011f98c0d5a93c5403094319226e1d77de0e6c7a Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Fri, 19 Dec 2025 14:43:59 +0530 Subject: [PATCH 2/9] feat:UploadFileAndFinishUploadSession --- docs/api_reference/file.rst | 1 + docs/getting_started.rst | 2 +- nisystemlink/clients/file/_file_client.py | 64 ++++++++++- .../models/_upload_session_start_response.py | 4 +- tests/integration/file/test_file_client.py | 101 ++++++++++++++++-- 5 files changed, 159 insertions(+), 13 deletions(-) diff --git a/docs/api_reference/file.rst b/docs/api_reference/file.rst index 4ef05fc2..5e72e930 100644 --- a/docs/api_reference/file.rst +++ b/docs/api_reference/file.rst @@ -16,6 +16,7 @@ nisystemlink.clients.file .. automethod:: download_file .. automethod:: update_metadata .. automethod:: start_upload_session + .. automethod:: append_to_upload_session .. automodule:: nisystemlink.clients.file.models :members: diff --git a/docs/getting_started.rst b/docs/getting_started.rst index bc6fc15e..fc47d369 100644 --- a/docs/getting_started.rst +++ b/docs/getting_started.rst @@ -204,7 +204,7 @@ default connection. The default connection depends on your environment. With a :class:`.FileClient` object, you can: * Get the list of files, download and delete files. -* Start upload sessions for chunked file uploads. +* Start upload sessions, upload file chunks, and finish sessions for large file uploads. Examples ~~~~~~~~ diff --git a/nisystemlink/clients/file/_file_client.py b/nisystemlink/clients/file/_file_client.py index c26a8a21..d23e26c7 100644 --- a/nisystemlink/clients/file/_file_client.py +++ b/nisystemlink/clients/file/_file_client.py @@ -276,7 +276,9 @@ def update_metadata(self, metadata: models.UpdateMetadataRequest, id: str) -> No ApiException: if unable to communicate with the File Service. """ - @post("service-groups/Default/upload-sessions/start", args=[Query(name="workspace")]) + @post( + "service-groups/Default/upload-sessions/start", args=[Query(name="workspace")] + ) def start_upload_session( self, workspace: str | None = None ) -> models.UploadSessionStartResponse: @@ -291,3 +293,63 @@ def start_upload_session( Raises: ApiException: if unable to communicate with the File Service. """ + + @post( + "service-groups/Default/upload-sessions/append", + args=[ + Query(name="sessionId"), + Query(name="chunk"), + Part(name="file"), + Query(name="close"), + ], + ) + @response_handler(lambda response: None) + def append_to_upload_session( + self, + session_id: str, + chunk: int, + file: BinaryIO, + close: bool | None = False, + ) -> None: + """Append a chunk to an upload session. + + The chunk needs to be 10485760 bytes (10 MB), unless the close parameter is true, + which means that it is the last chunk of the file content. The chunks can be uploaded + concurrently, as long as all chunk uploads are completed before finalizing. + + Args: + session_id: The id of the upload session. + chunk: The number of the chunk uploaded (0-based indexing). + file: The chunk data to upload. + close: Set the current chunk as the last chunk to be uploaded. Defaults to False. + + Raises: + ApiException: if unable to communicate with the File Service. + """ + + @response_handler(_file_uri_response_handler) + @post( + "service-groups/Default/upload-sessions/finish", + args=[Query(name="sessionId"), Field(name="name"), Field(name="properties")], + ) + def finish_upload_session( + self, + session_id: str, + name: str, + properties: Dict[str, str], + ) -> str: + """Finish an upload session and make the file visible in SystemLink. + + This will trigger file events, such as routines. + + Args: + session_id: The id of the upload session. + name: The name of the file. + properties: The properties of the file. + + Returns: + ID of the uploaded file. + + Raises: + ApiException: if unable to communicate with the File Service. + """ diff --git a/nisystemlink/clients/file/models/_upload_session_start_response.py b/nisystemlink/clients/file/models/_upload_session_start_response.py index 3d3396cc..a126fc42 100644 --- a/nisystemlink/clients/file/models/_upload_session_start_response.py +++ b/nisystemlink/clients/file/models/_upload_session_start_response.py @@ -9,13 +9,13 @@ class UploadSessionStartResponse(JsonModel): id: str """ The session id created. - + example: 54837669-8cf5-469e-bf7d-26cb808c8f24 """ created_at: datetime """ The date and time the upload session has started. - + example: 2018-05-15T18:54:27.519Z """ diff --git a/tests/integration/file/test_file_client.py b/tests/integration/file/test_file_client.py index 810374a2..b630aaf2 100644 --- a/tests/integration/file/test_file_client.py +++ b/tests/integration/file/test_file_client.py @@ -3,6 +3,7 @@ import io import string from datetime import datetime +from io import BytesIO from random import choices, randint from typing import BinaryIO @@ -264,15 +265,6 @@ def test__query_files_linq__filter_returns_no_results(self, client: FileClient): assert response.total_count.value == 0 assert response.total_count.relation == "eq" - def test__start_upload_session__returns_session_id(self, client: FileClient): - response = client.start_upload_session() - - assert response is not None - assert response.id is not None - assert isinstance(response.id, str) - assert response.created_at is not None - assert isinstance(response.created_at, datetime) - def test__start_upload_session__with_invalid_workspace__raises( self, client: FileClient ): @@ -280,3 +272,94 @@ def test__start_upload_session__with_invalid_workspace__raises( with pytest.raises(ApiException): client.start_upload_session(workspace=invalid_workspace_id) + + def test__append_to_upload_session__uploads_file_in_chunks( + self, client: FileClient + ): + # Create a test file with known content + test_content = b"A" * 10485760 + b"B" * 5000000 # 10 MB + 5 MB + chunk_size = 10485760 # 10 MB + + # Start upload session + session_response = client.start_upload_session() + + # Verify session response + assert session_response is not None + assert session_response.id is not None + assert isinstance(session_response.id, str) + assert session_response.created_at is not None + assert isinstance(session_response.created_at, datetime) + + session_id = session_response.id + file_id = None + + try: + # Upload first chunk + first_chunk = BytesIO(test_content[:chunk_size]) + client.append_to_upload_session( + session_id=session_id, chunk=1, file=first_chunk + ) + + # Upload second chunk (last chunk) + second_chunk = BytesIO(test_content[chunk_size:]) + client.append_to_upload_session( + session_id=session_id, chunk=2, file=second_chunk, close=True + ) + + # Finish the upload session + file_name = f"{PREFIX}chunked_upload_test.bin" + file_id = client.finish_upload_session( + session_id=session_id, + name=file_name, + properties={ + "Name": file_name, + "Test": "ChunkedUpload", + "Description": "Test file from chunked upload", + }, + ) + + # Verify the file was created with correct metadata + files = client.get_files(ids=[file_id]) + assert files.total_count == 1 + assert len(files.available_files) == 1 + assert files.available_files[0].id == file_id + assert files.available_files[0].properties is not None + assert files.available_files[0].properties.get("Name") == file_name + assert files.available_files[0].properties.get("Test") == "ChunkedUpload" + assert ( + files.available_files[0].properties.get("Description") + == "Test file from chunked upload" + ) + + # Verify file content + downloaded_data = client.download_file(id=file_id) + assert downloaded_data.read() == test_content + except ApiException as api_exception: + raise api_exception + # Finish the upload session if it failed during chunk upload + if not file_id: + file_name = f"{PREFIX}chunked_upload_test.bin" + file_id = client.finish_upload_session( + session_id=session_id, + name=file_name, + properties={"Name": file_name, "Test": "ChunkedUpload"}, + ) + raise api_exception + finally: + # Clean up + if file_id: + try: + client.delete_file(id=file_id) + except Exception: + pass # Ignore cleanup errors + + def test__finish_upload_session__invalid_session_id_raises( + self, client: FileClient, invalid_file_id: str + ): + file_name = f"{PREFIX}invalid_session.txt" + properties = {"Name": file_name} + + with pytest.raises(ApiException): + client.finish_upload_session( + session_id=invalid_file_id, name=file_name, properties=properties + ) From 60f92b9d0149deac1ea93426d3d039242fadf64b Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Mon, 22 Dec 2025 18:08:51 +0530 Subject: [PATCH 3/9] fix:PRComments --- examples/file/upload_file_chunked.py | 90 +++++++++++++++++++ nisystemlink/clients/file/_file_client.py | 6 +- .../models/_upload_session_start_response.py | 9 +- tests/integration/file/test_file_client.py | 12 ++- 4 files changed, 101 insertions(+), 16 deletions(-) create mode 100644 examples/file/upload_file_chunked.py diff --git a/examples/file/upload_file_chunked.py b/examples/file/upload_file_chunked.py new file mode 100644 index 00000000..08866ba7 --- /dev/null +++ b/examples/file/upload_file_chunked.py @@ -0,0 +1,90 @@ +"""Example to upload a large file to SystemLink using chunked upload (upload sessions). + +This example demonstrates how to upload a file in chunks using upload sessions. +This is useful for large files that need to be uploaded in multiple parts. +""" + +import io + +from nisystemlink.clients.core import HttpConfiguration +from nisystemlink.clients.file import FileClient + +# Configure connection to SystemLink server +server_configuration = HttpConfiguration( + server_uri="https://test-api.lifecyclesolutions.ni.com/", + api_key="zr7fUQj3R2zSBt6b46LGquPkPZJ8wll_wg6oqRLQn2", +) + +client = FileClient(configuration=server_configuration) + +# Generate example file content (20 MB for demonstration) +CHUNK_SIZE = 10 * 1024 * 1024 # 10 MB chunks +file_content = b"X" * (20 * 1024 * 1024) # 20 MB file + +# Step 1: Start an upload session +session_response = client.start_upload_session(workspace=None) +session_id = session_response.session_id +print(f"Started upload session with ID: {session_id}") + +# Step 2: Upload chunks +# Split the file content into chunks and upload them +file_id = None +try: + num_chunks = (len(file_content) + CHUNK_SIZE - 1) // CHUNK_SIZE + + for i in range(num_chunks): + start = i * CHUNK_SIZE + end = min(start + CHUNK_SIZE, len(file_content)) + chunk_data = file_content[start:end] + + # Create a file-like object for the chunk + chunk_file = io.BytesIO(chunk_data) + + # Determine if this is the last chunk + is_last_chunk = i == num_chunks - 1 + + # Upload the chunk (chunk_index is 0-based) + client.append_to_upload_session( + session_id=session_id, + chunk_index=i + 1, + file=chunk_file, + close=is_last_chunk, + ) + print(f"Uploaded chunk {i + 1}/{num_chunks} ({len(chunk_data)} bytes)") + + # Step 3: Finish the upload session + file_name = "large_file_example.bin" + properties = { + "Description": "Example file uploaded using chunked upload", + "FileSize": str(len(file_content)), + } + + file_id = client.finish_upload_session( + session_id=session_id, name=file_name, properties=properties + ) + + print(f"\nSuccessfully uploaded file '{file_name}' with FileID: {file_id}") + +except Exception as e: + print(f"Error during chunked upload: {e}") + # Attempt to finish the session to clean up resources on the server + try: + print("Attempting to clean up upload session...") + file_id = client.finish_upload_session( + session_id=session_id, + name=f"incomplete_{file_name}", + properties={"Status": "Incomplete", "Error": str(e)}, + ) + print(f"Session cleaned up. Partial file saved with FileID: {file_id}") + except Exception as cleanup_error: + print(f"Failed to clean up session: {cleanup_error}") + print(f"Upload session {session_id} may need manual cleanup") + +finally: + # Clean up: Delete the uploaded file (whether complete or incomplete) + if file_id: + try: + client.delete_file(id=file_id) + print(f"Deleted file (FileID: {file_id})") + except Exception as delete_error: + print(f"Failed to delete file: {delete_error}") diff --git a/nisystemlink/clients/file/_file_client.py b/nisystemlink/clients/file/_file_client.py index d23e26c7..34a5f209 100644 --- a/nisystemlink/clients/file/_file_client.py +++ b/nisystemlink/clients/file/_file_client.py @@ -307,9 +307,9 @@ def start_upload_session( def append_to_upload_session( self, session_id: str, - chunk: int, + chunk_index: int, file: BinaryIO, - close: bool | None = False, + close: bool = False, ) -> None: """Append a chunk to an upload session. @@ -319,7 +319,7 @@ def append_to_upload_session( Args: session_id: The id of the upload session. - chunk: The number of the chunk uploaded (0-based indexing). + chunk_index: The 0-based index of the chunk to be uploaded. file: The chunk data to upload. close: Set the current chunk as the last chunk to be uploaded. Defaults to False. diff --git a/nisystemlink/clients/file/models/_upload_session_start_response.py b/nisystemlink/clients/file/models/_upload_session_start_response.py index a126fc42..c87357c7 100644 --- a/nisystemlink/clients/file/models/_upload_session_start_response.py +++ b/nisystemlink/clients/file/models/_upload_session_start_response.py @@ -1,21 +1,18 @@ from datetime import datetime from nisystemlink.clients.core._uplink._json_model import JsonModel +from pydantic import Field class UploadSessionStartResponse(JsonModel): """Response model for starting an upload session.""" - id: str + session_id: str = Field(alias="id") """ - The session id created. - - example: 54837669-8cf5-469e-bf7d-26cb808c8f24 + The id created for the upload session. """ created_at: datetime """ The date and time the upload session has started. - - example: 2018-05-15T18:54:27.519Z """ diff --git a/tests/integration/file/test_file_client.py b/tests/integration/file/test_file_client.py index b630aaf2..b40c8f00 100644 --- a/tests/integration/file/test_file_client.py +++ b/tests/integration/file/test_file_client.py @@ -285,25 +285,25 @@ def test__append_to_upload_session__uploads_file_in_chunks( # Verify session response assert session_response is not None - assert session_response.id is not None - assert isinstance(session_response.id, str) + assert session_response.session_id is not None + assert isinstance(session_response.session_id, str) assert session_response.created_at is not None assert isinstance(session_response.created_at, datetime) - session_id = session_response.id + session_id = session_response.session_id file_id = None try: # Upload first chunk first_chunk = BytesIO(test_content[:chunk_size]) client.append_to_upload_session( - session_id=session_id, chunk=1, file=first_chunk + session_id=session_id, chunk_index=1, file=first_chunk ) # Upload second chunk (last chunk) second_chunk = BytesIO(test_content[chunk_size:]) client.append_to_upload_session( - session_id=session_id, chunk=2, file=second_chunk, close=True + session_id=session_id, chunk_index=2, file=second_chunk, close=True ) # Finish the upload session @@ -312,7 +312,6 @@ def test__append_to_upload_session__uploads_file_in_chunks( session_id=session_id, name=file_name, properties={ - "Name": file_name, "Test": "ChunkedUpload", "Description": "Test file from chunked upload", }, @@ -335,7 +334,6 @@ def test__append_to_upload_session__uploads_file_in_chunks( downloaded_data = client.download_file(id=file_id) assert downloaded_data.read() == test_content except ApiException as api_exception: - raise api_exception # Finish the upload session if it failed during chunk upload if not file_id: file_name = f"{PREFIX}chunked_upload_test.bin" From 89e47d1fb531b174c4d9922050d95306635483bd Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Mon, 22 Dec 2025 18:19:24 +0530 Subject: [PATCH 4/9] fix:ChunkIndexComments --- nisystemlink/clients/file/_file_client.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/nisystemlink/clients/file/_file_client.py b/nisystemlink/clients/file/_file_client.py index 34a5f209..c689a1b4 100644 --- a/nisystemlink/clients/file/_file_client.py +++ b/nisystemlink/clients/file/_file_client.py @@ -319,7 +319,7 @@ def append_to_upload_session( Args: session_id: The id of the upload session. - chunk_index: The 0-based index of the chunk to be uploaded. + chunk_index: The 1-based index of the chunk to be uploaded. file: The chunk data to upload. close: Set the current chunk as the last chunk to be uploaded. Defaults to False. From e706e47a8b42709f23d5216d0d4dc8db82860e6a Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Mon, 22 Dec 2025 18:38:12 +0530 Subject: [PATCH 5/9] fix:ExampleChanges --- examples/file/upload_file_chunked.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/examples/file/upload_file_chunked.py b/examples/file/upload_file_chunked.py index 08866ba7..8fe35fc9 100644 --- a/examples/file/upload_file_chunked.py +++ b/examples/file/upload_file_chunked.py @@ -11,8 +11,8 @@ # Configure connection to SystemLink server server_configuration = HttpConfiguration( - server_uri="https://test-api.lifecyclesolutions.ni.com/", - api_key="zr7fUQj3R2zSBt6b46LGquPkPZJ8wll_wg6oqRLQn2", + server_uri="https://yourserver.yourcompany.com", + api_key="YourAPIKeyGeneratedFromSystemLink", ) client = FileClient(configuration=server_configuration) From 1b88f104c9f62524482201ac4e2919b4e6c46c9d Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Mon, 22 Dec 2025 19:50:24 +0530 Subject: [PATCH 6/9] fix:ExampleFileChanges --- examples/file/upload_file_chunked.py | 229 ++++++++++++++++++++------- 1 file changed, 174 insertions(+), 55 deletions(-) diff --git a/examples/file/upload_file_chunked.py b/examples/file/upload_file_chunked.py index 8fe35fc9..37ad4da6 100644 --- a/examples/file/upload_file_chunked.py +++ b/examples/file/upload_file_chunked.py @@ -1,10 +1,16 @@ -"""Example to upload a large file to SystemLink using chunked upload (upload sessions). +"""Example comparing synchronous and asynchronous chunked file upload to SystemLink. -This example demonstrates how to upload a file in chunks using upload sessions. -This is useful for large files that need to be uploaded in multiple parts. +This example demonstrates: +1. Synchronous chunk upload (one chunk at a time) +2. Asynchronous chunk upload (multiple chunks concurrently) +3. Performance comparison between both approaches """ +import asyncio import io +import os +import tempfile +import time from nisystemlink.clients.core import HttpConfiguration from nisystemlink.clients.file import FileClient @@ -17,74 +23,187 @@ client = FileClient(configuration=server_configuration) -# Generate example file content (20 MB for demonstration) +# Generate example file content (50 MB for demonstration) CHUNK_SIZE = 10 * 1024 * 1024 # 10 MB chunks -file_content = b"X" * (20 * 1024 * 1024) # 20 MB file +FILE_SIZE = 50 * 1024 * 1024 # 50 MB file +# Generate test file content by repeating a simple message +test_data = b"This is test data for chunked file upload example.\n" +file_content = test_data * (FILE_SIZE // len(test_data)) + +# Create a temporary file to mimic file-on-disk behavior +with tempfile.NamedTemporaryFile(delete=False, suffix=".bin") as temp_file: + temp_file.write(file_content) + temp_file_path = temp_file.name + + +def upload_chunk( + session_id: str, chunk_index: int, chunk_data: bytes, is_last: bool +) -> int: + """Upload a single chunk.""" + chunk_file = io.BytesIO(chunk_data) + client.append_to_upload_session( + session_id=session_id, chunk_index=chunk_index, file=chunk_file, close=is_last + ) + return chunk_index -# Step 1: Start an upload session -session_response = client.start_upload_session(workspace=None) -session_id = session_response.session_id -print(f"Started upload session with ID: {session_id}") -# Step 2: Upload chunks -# Split the file content into chunks and upload them -file_id = None -try: - num_chunks = (len(file_content) + CHUNK_SIZE - 1) // CHUNK_SIZE +async def upload_chunk_async( + session_id: str, chunk_index: int, chunk_data: bytes, is_last: bool +) -> int: + """Upload a single chunk asynchronously.""" + # Run the synchronous upload in a thread pool to avoid blocking + return await asyncio.to_thread( + upload_chunk, session_id, chunk_index, chunk_data, is_last + ) + - for i in range(num_chunks): - start = i * CHUNK_SIZE - end = min(start + CHUNK_SIZE, len(file_content)) - chunk_data = file_content[start:end] +def upload_synchronous(): + """Upload file chunks synchronously (one at a time).""" + print("\nSynchronous Upload:") - # Create a file-like object for the chunk - chunk_file = io.BytesIO(chunk_data) + # Start upload session + session_response = client.start_upload_session(workspace=None) + session_id = session_response.session_id + print(f"Started upload session: {session_id}\n") + + file_id = None + start_time = time.time() + + try: + # Read and upload chunks sequentially using iter() with sentinel + with open(temp_file_path, "rb") as f: + chunks = list(enumerate(iter(lambda: f.read(CHUNK_SIZE), b""), start=1)) + num_chunks = len(chunks) - # Determine if this is the last chunk - is_last_chunk = i == num_chunks - 1 + for i, chunk_data in chunks: + is_last_chunk = i == num_chunks - # Upload the chunk (chunk_index is 0-based) - client.append_to_upload_session( + chunk_start = time.time() + upload_chunk(session_id, i, chunk_data, is_last_chunk) + chunk_time = time.time() - chunk_start + + print(f" Chunk {i}/{num_chunks} uploaded in {chunk_time:.2f}s") + + # Finish the upload session + file_id = client.finish_upload_session( session_id=session_id, - chunk_index=i + 1, - file=chunk_file, - close=is_last_chunk, + name="sync_upload_example.bin", + properties={"Type": "Synchronous", "FileSize": str(FILE_SIZE)}, ) - print(f"Uploaded chunk {i + 1}/{num_chunks} ({len(chunk_data)} bytes)") - # Step 3: Finish the upload session - file_name = "large_file_example.bin" - properties = { - "Description": "Example file uploaded using chunked upload", - "FileSize": str(len(file_content)), - } + total_time = time.time() - start_time + print(f"\nUpload completed in {total_time:.2f}s") + print(f"File ID: {file_id}\n") + + return file_id, total_time + + except Exception as e: + print(f"✗ Error: {e}") + if not file_id: + file_id = client.finish_upload_session( + session_id=session_id, + name="incomplete_sync.bin", + properties={"Status": "Failed"}, + ) + return file_id, None - file_id = client.finish_upload_session( - session_id=session_id, name=file_name, properties=properties - ) - print(f"\nSuccessfully uploaded file '{file_name}' with FileID: {file_id}") +async def upload_asynchronous(): + """Upload file chunks asynchronously (concurrently).""" + print("\nAsynchronous Upload (Concurrent):") + + # Start upload session + session_response = client.start_upload_session(workspace=None) + session_id = session_response.session_id + print(f"Started upload session: {session_id}\n") + + file_id = None + start_time = time.time() -except Exception as e: - print(f"Error during chunked upload: {e}") - # Attempt to finish the session to clean up resources on the server try: - print("Attempting to clean up upload session...") + # Read all chunks using iter() with sentinel for cleaner iteration + with open(temp_file_path, "rb") as f: + chunks = list(enumerate(iter(lambda: f.read(CHUNK_SIZE), b""), start=1)) + num_chunks = len(chunks) + chunks = [(i, data, i == num_chunks) for i, data in chunks] + + # Upload chunks concurrently + print(f"Uploading {num_chunks} chunks concurrently\n") + + tasks = [ + upload_chunk_async(session_id, idx, data, is_last) + for idx, data, is_last in chunks + ] + + # Run all upload tasks concurrently + results = await asyncio.gather(*tasks, return_exceptions=True) + + # Check for errors + for i, result in enumerate(results): + if isinstance(result, Exception): + print(f" Chunk {i + 1} failed: {result}") + raise result + else: + print(f" Chunk {result}/{num_chunks} uploaded") + + # Finish the upload session file_id = client.finish_upload_session( session_id=session_id, - name=f"incomplete_{file_name}", - properties={"Status": "Incomplete", "Error": str(e)}, + name="async_upload_example.bin", + properties={"Type": "Asynchronous", "FileSize": str(FILE_SIZE)}, ) - print(f"Session cleaned up. Partial file saved with FileID: {file_id}") - except Exception as cleanup_error: - print(f"Failed to clean up session: {cleanup_error}") - print(f"Upload session {session_id} may need manual cleanup") + + total_time = time.time() - start_time + print(f"\nUpload completed in {total_time:.2f}s") + print(f" File ID: {file_id}\n") + + return file_id, total_time + + except Exception as e: + print(f"✗ Error: {e}") + if not file_id: + file_id = client.finish_upload_session( + session_id=session_id, + name="incomplete_async.bin", + properties={"Status": "Failed"}, + ) + return file_id, None + + +# Run both upload methods and compare +async def main(): + """Main function to run both upload methods.""" + # Synchronous upload + sync_file_id, sync_time = upload_synchronous() + + # Asynchronous upload + async_file_id, async_time = await upload_asynchronous() + + # Performance comparison + if sync_time and async_time: + print("\nPerformance Comparison:") + print(f"Synchronous: {sync_time:.2f}s") + print(f"Asynchronous: {async_time:.2f}s") + + return sync_file_id, async_file_id + + +try: + # Run the async main function + sync_file_id, async_file_id = asyncio.run(main()) finally: - # Clean up: Delete the uploaded file (whether complete or incomplete) - if file_id: - try: - client.delete_file(id=file_id) - print(f"Deleted file (FileID: {file_id})") - except Exception as delete_error: - print(f"Failed to delete file: {delete_error}") + # Clean up: Delete uploaded files + print("Cleaning up...") + for file_id in [sync_file_id, async_file_id]: + if file_id: + try: + client.delete_file(id=file_id) + print(f" Deleted file: {file_id}") + except Exception as e: + print(f" Failed to delete {file_id}: {e}") + + os.unlink(temp_file_path) + print(f" Deleted temp file: {temp_file_path}") + + print("\nCleanup complete") From f7520f69ec3fdd4effb09ff6a7654ea77d190103 Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Tue, 23 Dec 2025 14:11:58 +0530 Subject: [PATCH 7/9] fix:PRComments --- docs/api_reference/file.rst | 1 + examples/file/upload_file_chunked.py | 209 --------------------- examples/file/upload_file_chunked_async.py | 143 ++++++++++++++ examples/file/upload_file_chunked_sync.py | 123 ++++++++++++ nisystemlink/clients/file/_file_client.py | 2 +- tests/integration/file/test_file_client.py | 4 +- 6 files changed, 270 insertions(+), 212 deletions(-) delete mode 100644 examples/file/upload_file_chunked.py create mode 100644 examples/file/upload_file_chunked_async.py create mode 100644 examples/file/upload_file_chunked_sync.py diff --git a/docs/api_reference/file.rst b/docs/api_reference/file.rst index 5e72e930..d21e17ff 100644 --- a/docs/api_reference/file.rst +++ b/docs/api_reference/file.rst @@ -17,6 +17,7 @@ nisystemlink.clients.file .. automethod:: update_metadata .. automethod:: start_upload_session .. automethod:: append_to_upload_session + .. automethod:: finish_upload_session .. automodule:: nisystemlink.clients.file.models :members: diff --git a/examples/file/upload_file_chunked.py b/examples/file/upload_file_chunked.py deleted file mode 100644 index 37ad4da6..00000000 --- a/examples/file/upload_file_chunked.py +++ /dev/null @@ -1,209 +0,0 @@ -"""Example comparing synchronous and asynchronous chunked file upload to SystemLink. - -This example demonstrates: -1. Synchronous chunk upload (one chunk at a time) -2. Asynchronous chunk upload (multiple chunks concurrently) -3. Performance comparison between both approaches -""" - -import asyncio -import io -import os -import tempfile -import time - -from nisystemlink.clients.core import HttpConfiguration -from nisystemlink.clients.file import FileClient - -# Configure connection to SystemLink server -server_configuration = HttpConfiguration( - server_uri="https://yourserver.yourcompany.com", - api_key="YourAPIKeyGeneratedFromSystemLink", -) - -client = FileClient(configuration=server_configuration) - -# Generate example file content (50 MB for demonstration) -CHUNK_SIZE = 10 * 1024 * 1024 # 10 MB chunks -FILE_SIZE = 50 * 1024 * 1024 # 50 MB file -# Generate test file content by repeating a simple message -test_data = b"This is test data for chunked file upload example.\n" -file_content = test_data * (FILE_SIZE // len(test_data)) - -# Create a temporary file to mimic file-on-disk behavior -with tempfile.NamedTemporaryFile(delete=False, suffix=".bin") as temp_file: - temp_file.write(file_content) - temp_file_path = temp_file.name - - -def upload_chunk( - session_id: str, chunk_index: int, chunk_data: bytes, is_last: bool -) -> int: - """Upload a single chunk.""" - chunk_file = io.BytesIO(chunk_data) - client.append_to_upload_session( - session_id=session_id, chunk_index=chunk_index, file=chunk_file, close=is_last - ) - return chunk_index - - -async def upload_chunk_async( - session_id: str, chunk_index: int, chunk_data: bytes, is_last: bool -) -> int: - """Upload a single chunk asynchronously.""" - # Run the synchronous upload in a thread pool to avoid blocking - return await asyncio.to_thread( - upload_chunk, session_id, chunk_index, chunk_data, is_last - ) - - -def upload_synchronous(): - """Upload file chunks synchronously (one at a time).""" - print("\nSynchronous Upload:") - - # Start upload session - session_response = client.start_upload_session(workspace=None) - session_id = session_response.session_id - print(f"Started upload session: {session_id}\n") - - file_id = None - start_time = time.time() - - try: - # Read and upload chunks sequentially using iter() with sentinel - with open(temp_file_path, "rb") as f: - chunks = list(enumerate(iter(lambda: f.read(CHUNK_SIZE), b""), start=1)) - num_chunks = len(chunks) - - for i, chunk_data in chunks: - is_last_chunk = i == num_chunks - - chunk_start = time.time() - upload_chunk(session_id, i, chunk_data, is_last_chunk) - chunk_time = time.time() - chunk_start - - print(f" Chunk {i}/{num_chunks} uploaded in {chunk_time:.2f}s") - - # Finish the upload session - file_id = client.finish_upload_session( - session_id=session_id, - name="sync_upload_example.bin", - properties={"Type": "Synchronous", "FileSize": str(FILE_SIZE)}, - ) - - total_time = time.time() - start_time - print(f"\nUpload completed in {total_time:.2f}s") - print(f"File ID: {file_id}\n") - - return file_id, total_time - - except Exception as e: - print(f"✗ Error: {e}") - if not file_id: - file_id = client.finish_upload_session( - session_id=session_id, - name="incomplete_sync.bin", - properties={"Status": "Failed"}, - ) - return file_id, None - - -async def upload_asynchronous(): - """Upload file chunks asynchronously (concurrently).""" - print("\nAsynchronous Upload (Concurrent):") - - # Start upload session - session_response = client.start_upload_session(workspace=None) - session_id = session_response.session_id - print(f"Started upload session: {session_id}\n") - - file_id = None - start_time = time.time() - - try: - # Read all chunks using iter() with sentinel for cleaner iteration - with open(temp_file_path, "rb") as f: - chunks = list(enumerate(iter(lambda: f.read(CHUNK_SIZE), b""), start=1)) - num_chunks = len(chunks) - chunks = [(i, data, i == num_chunks) for i, data in chunks] - - # Upload chunks concurrently - print(f"Uploading {num_chunks} chunks concurrently\n") - - tasks = [ - upload_chunk_async(session_id, idx, data, is_last) - for idx, data, is_last in chunks - ] - - # Run all upload tasks concurrently - results = await asyncio.gather(*tasks, return_exceptions=True) - - # Check for errors - for i, result in enumerate(results): - if isinstance(result, Exception): - print(f" Chunk {i + 1} failed: {result}") - raise result - else: - print(f" Chunk {result}/{num_chunks} uploaded") - - # Finish the upload session - file_id = client.finish_upload_session( - session_id=session_id, - name="async_upload_example.bin", - properties={"Type": "Asynchronous", "FileSize": str(FILE_SIZE)}, - ) - - total_time = time.time() - start_time - print(f"\nUpload completed in {total_time:.2f}s") - print(f" File ID: {file_id}\n") - - return file_id, total_time - - except Exception as e: - print(f"✗ Error: {e}") - if not file_id: - file_id = client.finish_upload_session( - session_id=session_id, - name="incomplete_async.bin", - properties={"Status": "Failed"}, - ) - return file_id, None - - -# Run both upload methods and compare -async def main(): - """Main function to run both upload methods.""" - # Synchronous upload - sync_file_id, sync_time = upload_synchronous() - - # Asynchronous upload - async_file_id, async_time = await upload_asynchronous() - - # Performance comparison - if sync_time and async_time: - print("\nPerformance Comparison:") - print(f"Synchronous: {sync_time:.2f}s") - print(f"Asynchronous: {async_time:.2f}s") - - return sync_file_id, async_file_id - - -try: - # Run the async main function - sync_file_id, async_file_id = asyncio.run(main()) - -finally: - # Clean up: Delete uploaded files - print("Cleaning up...") - for file_id in [sync_file_id, async_file_id]: - if file_id: - try: - client.delete_file(id=file_id) - print(f" Deleted file: {file_id}") - except Exception as e: - print(f" Failed to delete {file_id}: {e}") - - os.unlink(temp_file_path) - print(f" Deleted temp file: {temp_file_path}") - - print("\nCleanup complete") diff --git a/examples/file/upload_file_chunked_async.py b/examples/file/upload_file_chunked_async.py new file mode 100644 index 00000000..b7698b3f --- /dev/null +++ b/examples/file/upload_file_chunked_async.py @@ -0,0 +1,143 @@ +"""Example of asynchronous chunked file upload to SystemLink. + +This example demonstrates uploading a large file in chunks concurrently, +without loading the entire file into memory at once. Multiple chunks are +uploaded simultaneously for better performance. +""" + +import asyncio +import tempfile +import time +from functools import partial +from io import BytesIO + +from nisystemlink.clients.core import HttpConfiguration +from nisystemlink.clients.file import FileClient + +# Configure connection to SystemLink server +server_configuration = HttpConfiguration( + server_uri="https://yourserver.yourcompany.com", + api_key="YourAPIKeyGeneratedFromSystemLink", +) + +client = FileClient(configuration=server_configuration) + +# Generate example file content (50 MB for demonstration) +CHUNK_SIZE = 10 * 1024 * 1024 # 10 MB chunks +FILE_SIZE = 50 * 1024 * 1024 # 50 MB file +# Generate test file content by repeating a simple message +test_data = b"This is test data for chunked file upload example.\n" +file_content = test_data * (FILE_SIZE // len(test_data)) + + +async def upload_chunk_async( + session_id: str, chunk_index: int, chunk_data: bytes, is_last: bool +) -> int: + """Upload a single chunk asynchronously.""" + # Run the synchronous upload in a thread pool to avoid blocking + await asyncio.to_thread( + client.append_to_upload_session, + session_id=session_id, + chunk_index=chunk_index, + chunk=BytesIO(chunk_data), + close=is_last, + ) + return chunk_index + + +# Create a temporary file to mimic file-on-disk behavior +with tempfile.NamedTemporaryFile(delete=True, suffix=".bin") as temp_file: + temp_file.write(file_content) + temp_file_path = temp_file.name + temp_file.flush() # Ensure data is written to disk + + async def main(): + """Main async function to upload file chunks concurrently.""" + print(f"Created temporary file: {temp_file_path}") + print(f"File size: {FILE_SIZE / (1024 * 1024):.1f} MB") + print(f"Chunk size: {CHUNK_SIZE / (1024 * 1024):.1f} MB\n") + + # Start upload session + session_response = client.start_upload_session(workspace=None) + session_id = session_response.session_id + print(f"Started upload session: {session_id}\n") + + file_id = None + start_time = time.time() + + try: + chunks = [] + temp_file.seek(0) + + read_chunk = partial(temp_file.read, CHUNK_SIZE) + + chunk_iterator = iter(read_chunk, b"") + + chunk_index = 1 + for chunk_data in chunk_iterator: + chunks.append((chunk_index, chunk_data)) + chunk_index += 1 + + total_chunks = len(chunks) + print(f"Uploading {total_chunks} chunks concurrently...\n") + + # Upload chunks concurrently + tasks = [ + upload_chunk_async(session_id, idx, data, idx == total_chunks) + for idx, data in chunks + ] + + # Run all upload tasks concurrently + results = await asyncio.gather(*tasks, return_exceptions=True) + + # Check for errors + for i, result in enumerate(results): + if isinstance(result, Exception): + print(f"Chunk {i + 1} failed: {result}") + raise result + else: + print(f"Chunk {result}/{total_chunks} uploaded") + + # Finish the upload session + file_id = client.finish_upload_session( + session_id=session_id, + name="async_chunked_upload_example.bin", + properties={"Type": "Asynchronous Chunked", "FileSize": str(FILE_SIZE)}, + ) + + total_time = time.time() - start_time + print("\nUpload completed successfully!") + print(f"File ID: {file_id}") + print(f"Total time: {total_time:.2f}s") + + return file_id + + except Exception as e: + print(f"Error during upload: {e}") + if not file_id: + try: + file_id = client.finish_upload_session( + session_id=session_id, + name="incomplete_async_upload.bin", + properties={"Status": "Failed"}, + ) + except Exception as finish_error: + print(f"Failed to finish session: {finish_error}") + return file_id + + # Run the async main function + try: + file_id = asyncio.run(main()) + + finally: + # Clean up: Delete uploaded file + print("\nCleaning up...") + if file_id: + try: + client.delete_file(id=file_id) + print(f"Deleted uploaded file: {file_id}") + except Exception as e: + print(f"Failed to delete file {file_id}: {e}") + + print(f"Temporary file will be automatically deleted: {temp_file_path}") + print("Cleanup complete") diff --git a/examples/file/upload_file_chunked_sync.py b/examples/file/upload_file_chunked_sync.py new file mode 100644 index 00000000..a2c81ea8 --- /dev/null +++ b/examples/file/upload_file_chunked_sync.py @@ -0,0 +1,123 @@ +"""Example of synchronous chunked file upload to SystemLink. + +This example demonstrates uploading a large file in chunks without loading +the entire file into memory at once. Each chunk is read and uploaded sequentially. +""" + +import tempfile +import time +from functools import partial +from io import BytesIO + +from nisystemlink.clients.core import HttpConfiguration +from nisystemlink.clients.file import FileClient + +# Configure connection to SystemLink server +server_configuration = HttpConfiguration( + server_uri="https://yourserver.yourcompany.com", + api_key="YourAPIKeyGeneratedFromSystemLink", +) + +client = FileClient(configuration=server_configuration) + +# Generate example file content (50 MB for demonstration) +CHUNK_SIZE = 10 * 1024 * 1024 # 10 MB chunks +FILE_SIZE = 50 * 1024 * 1024 # 50 MB file +# Generate test file content by repeating a simple message +test_data = b"This is test data for chunked file upload example.\n" +file_content = test_data * (FILE_SIZE // len(test_data)) + +# Create a temporary file to mimic file-on-disk behavior +# delete=True ensures automatic cleanup when the with block exits +with tempfile.NamedTemporaryFile(delete=True, suffix=".bin") as temp_file: + temp_file.write(file_content) + temp_file_path = temp_file.name + temp_file.flush() # Ensure data is written to disk + + print(f"Created temporary file: {temp_file_path}") + print(f"File size: {FILE_SIZE / (1024 * 1024):.1f} MB") + print(f"Chunk size: {CHUNK_SIZE / (1024 * 1024):.1f} MB\n") + + # Start upload session + session_response = client.start_upload_session(workspace=None) + session_id = session_response.session_id + print(f"Started upload session: {session_id}\n") + + file_id = None + start_time = time.time() + + try: + # Read and upload chunks sequentially using iter() with sentinel + # This approach only loads one chunk into memory at a time + temp_file.seek(0) # Seek back to the beginning of the file + + # Create a partial function for reading chunks + read_chunk = partial(temp_file.read, CHUNK_SIZE) + + # Create an iterator that yields chunks until an empty bytes object is returned + chunk_iterator = iter(read_chunk, b"") + + # Process chunks one at a time + chunk_index = 1 + total_chunks = ( + FILE_SIZE + CHUNK_SIZE - 1 + ) // CHUNK_SIZE # Calculate total chunks + + for chunk_data in chunk_iterator: + is_last_chunk = len(chunk_data) < CHUNK_SIZE or chunk_index == total_chunks + + chunk_start = time.time() + + # Upload chunk directly (wrap in BytesIO for type compatibility) + client.append_to_upload_session( + session_id=session_id, + chunk_index=chunk_index, + chunk=BytesIO(chunk_data), + close=is_last_chunk, + ) + + chunk_time = time.time() - chunk_start + chunk_size_mb = len(chunk_data) / (1024 * 1024) + print( + f"Chunk {chunk_index}/{total_chunks} uploaded " + f"({chunk_size_mb:.1f} MB) in {chunk_time:.2f}s" + ) + + chunk_index += 1 + + # Finish the upload session + file_id = client.finish_upload_session( + session_id=session_id, + name="chunked_upload_example.bin", + properties={"Type": "Synchronous Chunked", "FileSize": str(FILE_SIZE)}, + ) + + total_time = time.time() - start_time + print("\nUpload completed successfully!") + print(f"File ID: {file_id}") + print(f"Total time: {total_time:.2f}s") + + except Exception as e: + print(f"Error during upload: {e}") + if not file_id: + try: + file_id = client.finish_upload_session( + session_id=session_id, + name="incomplete_upload.bin", + properties={"Status": "Failed"}, + ) + except Exception as finish_error: + print(f"Failed to finish session: {finish_error}") + + finally: + # Clean up: Delete uploaded file + print("\nCleaning up...") + if file_id: + try: + client.delete_file(id=file_id) + print(f"Deleted uploaded file: {file_id}") + except Exception as e: + print(f"Failed to delete file {file_id}: {e}") + + print(f"Temporary file will be automatically deleted: {temp_file_path}") + print("Cleanup complete") diff --git a/nisystemlink/clients/file/_file_client.py b/nisystemlink/clients/file/_file_client.py index c689a1b4..852f24b5 100644 --- a/nisystemlink/clients/file/_file_client.py +++ b/nisystemlink/clients/file/_file_client.py @@ -308,7 +308,7 @@ def append_to_upload_session( self, session_id: str, chunk_index: int, - file: BinaryIO, + chunk: BinaryIO, close: bool = False, ) -> None: """Append a chunk to an upload session. diff --git a/tests/integration/file/test_file_client.py b/tests/integration/file/test_file_client.py index b40c8f00..3f9d0912 100644 --- a/tests/integration/file/test_file_client.py +++ b/tests/integration/file/test_file_client.py @@ -297,13 +297,13 @@ def test__append_to_upload_session__uploads_file_in_chunks( # Upload first chunk first_chunk = BytesIO(test_content[:chunk_size]) client.append_to_upload_session( - session_id=session_id, chunk_index=1, file=first_chunk + session_id=session_id, chunk_index=1, chunk=first_chunk ) # Upload second chunk (last chunk) second_chunk = BytesIO(test_content[chunk_size:]) client.append_to_upload_session( - session_id=session_id, chunk_index=2, file=second_chunk, close=True + session_id=session_id, chunk_index=2, chunk=second_chunk, close=True ) # Finish the upload session From f49b89b45b27089680a32bb2eba70f6c836be098 Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Wed, 24 Dec 2025 17:10:04 +0530 Subject: [PATCH 8/9] fix:Remove Example Files --- examples/file/upload_file_chunked_async.py | 143 --------------------- examples/file/upload_file_chunked_sync.py | 123 ------------------ 2 files changed, 266 deletions(-) delete mode 100644 examples/file/upload_file_chunked_async.py delete mode 100644 examples/file/upload_file_chunked_sync.py diff --git a/examples/file/upload_file_chunked_async.py b/examples/file/upload_file_chunked_async.py deleted file mode 100644 index b7698b3f..00000000 --- a/examples/file/upload_file_chunked_async.py +++ /dev/null @@ -1,143 +0,0 @@ -"""Example of asynchronous chunked file upload to SystemLink. - -This example demonstrates uploading a large file in chunks concurrently, -without loading the entire file into memory at once. Multiple chunks are -uploaded simultaneously for better performance. -""" - -import asyncio -import tempfile -import time -from functools import partial -from io import BytesIO - -from nisystemlink.clients.core import HttpConfiguration -from nisystemlink.clients.file import FileClient - -# Configure connection to SystemLink server -server_configuration = HttpConfiguration( - server_uri="https://yourserver.yourcompany.com", - api_key="YourAPIKeyGeneratedFromSystemLink", -) - -client = FileClient(configuration=server_configuration) - -# Generate example file content (50 MB for demonstration) -CHUNK_SIZE = 10 * 1024 * 1024 # 10 MB chunks -FILE_SIZE = 50 * 1024 * 1024 # 50 MB file -# Generate test file content by repeating a simple message -test_data = b"This is test data for chunked file upload example.\n" -file_content = test_data * (FILE_SIZE // len(test_data)) - - -async def upload_chunk_async( - session_id: str, chunk_index: int, chunk_data: bytes, is_last: bool -) -> int: - """Upload a single chunk asynchronously.""" - # Run the synchronous upload in a thread pool to avoid blocking - await asyncio.to_thread( - client.append_to_upload_session, - session_id=session_id, - chunk_index=chunk_index, - chunk=BytesIO(chunk_data), - close=is_last, - ) - return chunk_index - - -# Create a temporary file to mimic file-on-disk behavior -with tempfile.NamedTemporaryFile(delete=True, suffix=".bin") as temp_file: - temp_file.write(file_content) - temp_file_path = temp_file.name - temp_file.flush() # Ensure data is written to disk - - async def main(): - """Main async function to upload file chunks concurrently.""" - print(f"Created temporary file: {temp_file_path}") - print(f"File size: {FILE_SIZE / (1024 * 1024):.1f} MB") - print(f"Chunk size: {CHUNK_SIZE / (1024 * 1024):.1f} MB\n") - - # Start upload session - session_response = client.start_upload_session(workspace=None) - session_id = session_response.session_id - print(f"Started upload session: {session_id}\n") - - file_id = None - start_time = time.time() - - try: - chunks = [] - temp_file.seek(0) - - read_chunk = partial(temp_file.read, CHUNK_SIZE) - - chunk_iterator = iter(read_chunk, b"") - - chunk_index = 1 - for chunk_data in chunk_iterator: - chunks.append((chunk_index, chunk_data)) - chunk_index += 1 - - total_chunks = len(chunks) - print(f"Uploading {total_chunks} chunks concurrently...\n") - - # Upload chunks concurrently - tasks = [ - upload_chunk_async(session_id, idx, data, idx == total_chunks) - for idx, data in chunks - ] - - # Run all upload tasks concurrently - results = await asyncio.gather(*tasks, return_exceptions=True) - - # Check for errors - for i, result in enumerate(results): - if isinstance(result, Exception): - print(f"Chunk {i + 1} failed: {result}") - raise result - else: - print(f"Chunk {result}/{total_chunks} uploaded") - - # Finish the upload session - file_id = client.finish_upload_session( - session_id=session_id, - name="async_chunked_upload_example.bin", - properties={"Type": "Asynchronous Chunked", "FileSize": str(FILE_SIZE)}, - ) - - total_time = time.time() - start_time - print("\nUpload completed successfully!") - print(f"File ID: {file_id}") - print(f"Total time: {total_time:.2f}s") - - return file_id - - except Exception as e: - print(f"Error during upload: {e}") - if not file_id: - try: - file_id = client.finish_upload_session( - session_id=session_id, - name="incomplete_async_upload.bin", - properties={"Status": "Failed"}, - ) - except Exception as finish_error: - print(f"Failed to finish session: {finish_error}") - return file_id - - # Run the async main function - try: - file_id = asyncio.run(main()) - - finally: - # Clean up: Delete uploaded file - print("\nCleaning up...") - if file_id: - try: - client.delete_file(id=file_id) - print(f"Deleted uploaded file: {file_id}") - except Exception as e: - print(f"Failed to delete file {file_id}: {e}") - - print(f"Temporary file will be automatically deleted: {temp_file_path}") - print("Cleanup complete") diff --git a/examples/file/upload_file_chunked_sync.py b/examples/file/upload_file_chunked_sync.py deleted file mode 100644 index a2c81ea8..00000000 --- a/examples/file/upload_file_chunked_sync.py +++ /dev/null @@ -1,123 +0,0 @@ -"""Example of synchronous chunked file upload to SystemLink. - -This example demonstrates uploading a large file in chunks without loading -the entire file into memory at once. Each chunk is read and uploaded sequentially. -""" - -import tempfile -import time -from functools import partial -from io import BytesIO - -from nisystemlink.clients.core import HttpConfiguration -from nisystemlink.clients.file import FileClient - -# Configure connection to SystemLink server -server_configuration = HttpConfiguration( - server_uri="https://yourserver.yourcompany.com", - api_key="YourAPIKeyGeneratedFromSystemLink", -) - -client = FileClient(configuration=server_configuration) - -# Generate example file content (50 MB for demonstration) -CHUNK_SIZE = 10 * 1024 * 1024 # 10 MB chunks -FILE_SIZE = 50 * 1024 * 1024 # 50 MB file -# Generate test file content by repeating a simple message -test_data = b"This is test data for chunked file upload example.\n" -file_content = test_data * (FILE_SIZE // len(test_data)) - -# Create a temporary file to mimic file-on-disk behavior -# delete=True ensures automatic cleanup when the with block exits -with tempfile.NamedTemporaryFile(delete=True, suffix=".bin") as temp_file: - temp_file.write(file_content) - temp_file_path = temp_file.name - temp_file.flush() # Ensure data is written to disk - - print(f"Created temporary file: {temp_file_path}") - print(f"File size: {FILE_SIZE / (1024 * 1024):.1f} MB") - print(f"Chunk size: {CHUNK_SIZE / (1024 * 1024):.1f} MB\n") - - # Start upload session - session_response = client.start_upload_session(workspace=None) - session_id = session_response.session_id - print(f"Started upload session: {session_id}\n") - - file_id = None - start_time = time.time() - - try: - # Read and upload chunks sequentially using iter() with sentinel - # This approach only loads one chunk into memory at a time - temp_file.seek(0) # Seek back to the beginning of the file - - # Create a partial function for reading chunks - read_chunk = partial(temp_file.read, CHUNK_SIZE) - - # Create an iterator that yields chunks until an empty bytes object is returned - chunk_iterator = iter(read_chunk, b"") - - # Process chunks one at a time - chunk_index = 1 - total_chunks = ( - FILE_SIZE + CHUNK_SIZE - 1 - ) // CHUNK_SIZE # Calculate total chunks - - for chunk_data in chunk_iterator: - is_last_chunk = len(chunk_data) < CHUNK_SIZE or chunk_index == total_chunks - - chunk_start = time.time() - - # Upload chunk directly (wrap in BytesIO for type compatibility) - client.append_to_upload_session( - session_id=session_id, - chunk_index=chunk_index, - chunk=BytesIO(chunk_data), - close=is_last_chunk, - ) - - chunk_time = time.time() - chunk_start - chunk_size_mb = len(chunk_data) / (1024 * 1024) - print( - f"Chunk {chunk_index}/{total_chunks} uploaded " - f"({chunk_size_mb:.1f} MB) in {chunk_time:.2f}s" - ) - - chunk_index += 1 - - # Finish the upload session - file_id = client.finish_upload_session( - session_id=session_id, - name="chunked_upload_example.bin", - properties={"Type": "Synchronous Chunked", "FileSize": str(FILE_SIZE)}, - ) - - total_time = time.time() - start_time - print("\nUpload completed successfully!") - print(f"File ID: {file_id}") - print(f"Total time: {total_time:.2f}s") - - except Exception as e: - print(f"Error during upload: {e}") - if not file_id: - try: - file_id = client.finish_upload_session( - session_id=session_id, - name="incomplete_upload.bin", - properties={"Status": "Failed"}, - ) - except Exception as finish_error: - print(f"Failed to finish session: {finish_error}") - - finally: - # Clean up: Delete uploaded file - print("\nCleaning up...") - if file_id: - try: - client.delete_file(id=file_id) - print(f"Deleted uploaded file: {file_id}") - except Exception as e: - print(f"Failed to delete file {file_id}: {e}") - - print(f"Temporary file will be automatically deleted: {temp_file_path}") - print("Cleanup complete") From 2947fe750fb276d7c79625f223bb3340ddaa59e9 Mon Sep 17 00:00:00 2001 From: RSam-NI Date: Mon, 5 Jan 2026 12:38:28 +0530 Subject: [PATCH 9/9] fix:PRComments --- tests/integration/file/test_file_client.py | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/tests/integration/file/test_file_client.py b/tests/integration/file/test_file_client.py index 3f9d0912..234c0d79 100644 --- a/tests/integration/file/test_file_client.py +++ b/tests/integration/file/test_file_client.py @@ -333,16 +333,20 @@ def test__append_to_upload_session__uploads_file_in_chunks( # Verify file content downloaded_data = client.download_file(id=file_id) assert downloaded_data.read() == test_content - except ApiException as api_exception: + except ApiException: # Finish the upload session if it failed during chunk upload if not file_id: file_name = f"{PREFIX}chunked_upload_test.bin" - file_id = client.finish_upload_session( - session_id=session_id, - name=file_name, - properties={"Name": file_name, "Test": "ChunkedUpload"}, - ) - raise api_exception + try: + file_id = client.finish_upload_session( + session_id=session_id, + name=file_name, + properties={"Name": file_name, "Test": "ChunkedUpload"}, + ) + except ApiException: + pass + + raise finally: # Clean up if file_id: