From d858616319785c9b726d346978835c522fb61506 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 4 Mar 2026 15:03:46 +0000 Subject: [PATCH] BlobStore (trustgraph-flow/trustgraph/librarian/blob_store.py): - get_stream() - yields document content in chunks for streaming retrieval - create_multipart_upload() - initializes S3 multipart upload, returns upload_id - upload_part() - uploads a single part, returns etag - complete_multipart_upload() - finalizes upload with part etags - abort_multipart_upload() - cancels and cleans up Cassandra schema (trustgraph-flow/trustgraph/tables/library.py): - New upload_session table with 24-hour TTL - Index on user for listing sessions - Prepared statements for all operations - Methods: create_upload_session(), get_upload_session(), update_upload_session_chunk(), delete_upload_session(), list_upload_sessions() --- .../trustgraph/librarian/blob_store.py | 141 ++++++++++++ trustgraph-flow/trustgraph/tables/library.py | 216 ++++++++++++++++++ 2 files changed, 357 insertions(+) diff --git a/trustgraph-flow/trustgraph/librarian/blob_store.py b/trustgraph-flow/trustgraph/librarian/blob_store.py index 436e2718..e7949584 100644 --- a/trustgraph-flow/trustgraph/librarian/blob_store.py +++ b/trustgraph-flow/trustgraph/librarian/blob_store.py @@ -3,9 +3,12 @@ from .. knowledge import hash from .. exceptions import RequestError from minio import Minio +from minio.commonconfig import Part import time import io import logging +from typing import Iterator, List, Tuple +from uuid import UUID # Module logger logger = logging.getLogger(__name__) @@ -78,3 +81,141 @@ class BlobStore: return resp.read() + def get_stream(self, object_id, chunk_size: int = 1024 * 1024) -> Iterator[bytes]: + """ + Stream document content in chunks. + + Yields chunks of the document, allowing processing without loading + the entire document into memory. + + Args: + object_id: The UUID of the document object + chunk_size: Size of each chunk in bytes (default 1MB) + + Yields: + Chunks of document content as bytes + """ + resp = self.client.get_object( + bucket_name=self.bucket_name, + object_name="doc/" + str(object_id), + ) + + try: + while True: + chunk = resp.read(chunk_size) + if not chunk: + break + yield chunk + finally: + resp.close() + resp.release_conn() + + logger.debug("Stream complete") + + def create_multipart_upload(self, object_id: UUID, kind: str) -> str: + """ + Initialize a multipart upload. + + Args: + object_id: The UUID for the new object + kind: MIME type of the document + + Returns: + The S3 upload_id for this multipart upload session + """ + object_name = "doc/" + str(object_id) + + # Use minio's internal method to create multipart upload + upload_id = self.client._create_multipart_upload( + bucket_name=self.bucket_name, + object_name=object_name, + headers={"Content-Type": kind}, + ) + + logger.info(f"Created multipart upload {upload_id} for {object_id}") + return upload_id + + def upload_part( + self, + object_id: UUID, + upload_id: str, + part_number: int, + data: bytes + ) -> str: + """ + Upload a single part of a multipart upload. + + Args: + object_id: The UUID of the object being uploaded + upload_id: The S3 upload_id from create_multipart_upload + part_number: Part number (1-indexed, as per S3 spec) + data: The chunk data to upload + + Returns: + The ETag for this part (needed for complete_multipart_upload) + """ + object_name = "doc/" + str(object_id) + + etag = self.client._upload_part( + bucket_name=self.bucket_name, + object_name=object_name, + upload_id=upload_id, + part_number=part_number, + data=io.BytesIO(data), + headers={"Content-Length": str(len(data))}, + ) + + logger.debug(f"Uploaded part {part_number} for {object_id}, etag={etag}") + return etag + + def complete_multipart_upload( + self, + object_id: UUID, + upload_id: str, + parts: List[Tuple[int, str]] + ) -> None: + """ + Complete a multipart upload, assembling all parts into the final object. + + S3 coalesces the parts server-side - no data transfer through this client. + + Args: + object_id: The UUID of the object + upload_id: The S3 upload_id from create_multipart_upload + parts: List of (part_number, etag) tuples in order + """ + object_name = "doc/" + str(object_id) + + # Convert to Part objects as expected by minio + part_objects = [ + Part(part_number, etag) + for part_number, etag in parts + ] + + self.client._complete_multipart_upload( + bucket_name=self.bucket_name, + object_name=object_name, + upload_id=upload_id, + parts=part_objects, + ) + + logger.info(f"Completed multipart upload for {object_id}") + + def abort_multipart_upload(self, object_id: UUID, upload_id: str) -> None: + """ + Abort a multipart upload, cleaning up any uploaded parts. + + Args: + object_id: The UUID of the object + upload_id: The S3 upload_id from create_multipart_upload + """ + object_name = "doc/" + str(object_id) + + self.client._abort_multipart_upload( + bucket_name=self.bucket_name, + object_name=object_name, + upload_id=upload_id, + ) + + logger.info(f"Aborted multipart upload {upload_id} for {object_id}") + diff --git a/trustgraph-flow/trustgraph/tables/library.py b/trustgraph-flow/trustgraph/tables/library.py index 8bbe2bad..aa73c950 100644 --- a/trustgraph-flow/trustgraph/tables/library.py +++ b/trustgraph-flow/trustgraph/tables/library.py @@ -127,6 +127,32 @@ class LibraryTableStore: ); """); + logger.debug("upload_session table...") + + self.cassandra.execute(""" + CREATE TABLE IF NOT EXISTS upload_session ( + upload_id text PRIMARY KEY, + user text, + document_id text, + document_metadata text, + s3_upload_id text, + object_id uuid, + total_size bigint, + chunk_size int, + total_chunks int, + chunks_received map, + created_at timestamp, + updated_at timestamp + ) WITH default_time_to_live = 86400; + """); + + logger.debug("upload_session user index...") + + self.cassandra.execute(""" + CREATE INDEX IF NOT EXISTS upload_session_user + ON upload_session (user) + """); + logger.info("Cassandra schema OK.") def prepare_statements(self): @@ -210,6 +236,47 @@ class LibraryTableStore: WHERE user = ? """) + # Upload session prepared statements + self.insert_upload_session_stmt = self.cassandra.prepare(""" + INSERT INTO upload_session + ( + upload_id, user, document_id, document_metadata, + s3_upload_id, object_id, total_size, chunk_size, + total_chunks, chunks_received, created_at, updated_at + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """) + + self.get_upload_session_stmt = self.cassandra.prepare(""" + SELECT + upload_id, user, document_id, document_metadata, + s3_upload_id, object_id, total_size, chunk_size, + total_chunks, chunks_received, created_at, updated_at + FROM upload_session + WHERE upload_id = ? + """) + + self.update_upload_session_chunk_stmt = self.cassandra.prepare(""" + UPDATE upload_session + SET chunks_received = chunks_received + ?, + updated_at = ? + WHERE upload_id = ? + """) + + self.delete_upload_session_stmt = self.cassandra.prepare(""" + DELETE FROM upload_session + WHERE upload_id = ? + """) + + self.list_upload_sessions_stmt = self.cassandra.prepare(""" + SELECT + upload_id, document_id, document_metadata, + total_size, chunk_size, total_chunks, + chunks_received, created_at, updated_at + FROM upload_session + WHERE user = ? + """) + async def document_exists(self, user, id): resp = self.cassandra.execute( @@ -532,3 +599,152 @@ class LibraryTableStore: logger.debug("Done") return lst + + # Upload session methods + + async def create_upload_session( + self, + upload_id, + user, + document_id, + document_metadata, + s3_upload_id, + object_id, + total_size, + chunk_size, + total_chunks, + ): + """Create a new upload session for chunked upload.""" + + logger.info(f"Creating upload session {upload_id}") + + now = int(time.time() * 1000) + + while True: + try: + self.cassandra.execute( + self.insert_upload_session_stmt, + ( + upload_id, user, document_id, document_metadata, + s3_upload_id, object_id, total_size, chunk_size, + total_chunks, {}, now, now + ) + ) + break + except Exception as e: + logger.error("Exception occurred", exc_info=True) + raise e + + logger.debug("Upload session created") + + async def get_upload_session(self, upload_id): + """Get an upload session by ID.""" + + logger.debug(f"Get upload session {upload_id}") + + while True: + try: + resp = self.cassandra.execute( + self.get_upload_session_stmt, + (upload_id,) + ) + break + except Exception as e: + logger.error("Exception occurred", exc_info=True) + raise e + + for row in resp: + session = { + "upload_id": row[0], + "user": row[1], + "document_id": row[2], + "document_metadata": row[3], + "s3_upload_id": row[4], + "object_id": row[5], + "total_size": row[6], + "chunk_size": row[7], + "total_chunks": row[8], + "chunks_received": row[9] if row[9] else {}, + "created_at": row[10], + "updated_at": row[11], + } + logger.debug("Done") + return session + + return None + + async def update_upload_session_chunk(self, upload_id, chunk_index, etag): + """Record a successfully uploaded chunk.""" + + logger.debug(f"Update upload session {upload_id} chunk {chunk_index}") + + now = int(time.time() * 1000) + + while True: + try: + self.cassandra.execute( + self.update_upload_session_chunk_stmt, + ( + {chunk_index: etag}, + now, + upload_id + ) + ) + break + except Exception as e: + logger.error("Exception occurred", exc_info=True) + raise e + + logger.debug("Chunk recorded") + + async def delete_upload_session(self, upload_id): + """Delete an upload session.""" + + logger.info(f"Deleting upload session {upload_id}") + + while True: + try: + self.cassandra.execute( + self.delete_upload_session_stmt, + (upload_id,) + ) + break + except Exception as e: + logger.error("Exception occurred", exc_info=True) + raise e + + logger.debug("Upload session deleted") + + async def list_upload_sessions(self, user): + """List all upload sessions for a user.""" + + logger.debug(f"List upload sessions for {user}") + + while True: + try: + resp = self.cassandra.execute( + self.list_upload_sessions_stmt, + (user,) + ) + break + except Exception as e: + logger.error("Exception occurred", exc_info=True) + raise e + + sessions = [] + for row in resp: + chunks_received = row[6] if row[6] else {} + sessions.append({ + "upload_id": row[0], + "document_id": row[1], + "document_metadata": row[2], + "total_size": row[3], + "chunk_size": row[4], + "total_chunks": row[5], + "chunks_received": len(chunks_received), + "created_at": row[7], + "updated_at": row[8], + }) + + logger.debug("Done") + return sessions