diff --git a/trustgraph-flow/trustgraph/direct/cassandra_kg.py b/trustgraph-flow/trustgraph/direct/cassandra_kg.py index 711aef83..f44c8975 100644 --- a/trustgraph-flow/trustgraph/direct/cassandra_kg.py +++ b/trustgraph-flow/trustgraph/direct/cassandra_kg.py @@ -247,6 +247,32 @@ class KnowledgeGraph: (collection, s, p, o, limit) ) + def collection_exists(self, collection): + """Check if collection exists by querying triples_collection table""" + try: + result = self.session.execute( + f"SELECT collection FROM {self.collection_table} WHERE collection = %s LIMIT 1", + (collection,) + ) + return bool(list(result)) + except Exception as e: + logger.error(f"Error checking collection existence: {e}") + return False + + def create_collection(self, collection): + """Create collection by inserting marker row in triples_collection table""" + try: + # Insert a system marker triple to establish collection exists + # This won't interfere with data as it uses a reserved system namespace + self.session.execute( + f"INSERT INTO {self.collection_table} (collection, s, p, o) VALUES (%s, %s, %s, %s)", + (collection, "__system__", "__collection_created__", "__marker__") + ) + logger.info(f"Created collection marker for {collection}") + except Exception as e: + logger.error(f"Error creating collection: {e}") + raise e + def delete_collection(self, collection): """Delete all triples for a specific collection diff --git a/trustgraph-flow/trustgraph/librarian/collection_manager.py b/trustgraph-flow/trustgraph/librarian/collection_manager.py index f830db28..ea67fa3f 100644 --- a/trustgraph-flow/trustgraph/librarian/collection_manager.py +++ b/trustgraph-flow/trustgraph/librarian/collection_manager.py @@ -60,7 +60,7 @@ class CollectionManager: async def ensure_collection_exists(self, user: str, collection: str): """ - Ensure a collection exists, creating it if necessary (lazy creation) + Ensure a collection exists, creating it if necessary with broadcast to storage Args: user: User ID @@ -74,7 +74,7 @@ class CollectionManager: return # Create new collection with default metadata - logger.info(f"Creating new collection {user}/{collection}") + logger.info(f"Auto-creating collection {user}/{collection} from document submission") await self.table_store.create_collection( user=user, collection=collection, @@ -83,10 +83,64 @@ class CollectionManager: tags=set() ) + # Broadcast collection creation to all storage backends + creation_key = (user, collection) + logger.info(f"Broadcasting create-collection for {creation_key}") + + self.pending_deletions[creation_key] = { + "responses_pending": 3, # vector, object, triples + "responses_received": [], + "all_successful": True, + "error_messages": [], + "deletion_complete": asyncio.Event() + } + + storage_request = StorageManagementRequest( + operation="create-collection", + user=user, + collection=collection + ) + + # Send creation requests to all storage types + if self.vector_storage_producer: + await self.vector_storage_producer.send(storage_request) + if self.object_storage_producer: + await self.object_storage_producer.send(storage_request) + if self.triples_storage_producer: + await self.triples_storage_producer.send(storage_request) + + # Wait for all storage creations to complete (with timeout) + creation_info = self.pending_deletions[creation_key] + try: + await asyncio.wait_for( + creation_info["deletion_complete"].wait(), + timeout=30.0 # 30 second timeout + ) + except asyncio.TimeoutError: + logger.error(f"Timeout waiting for storage creation responses for {creation_key}") + creation_info["all_successful"] = False + creation_info["error_messages"].append("Timeout waiting for storage creation") + + # Check if all creations succeeded + if not creation_info["all_successful"]: + error_msg = f"Storage creation failed: {'; '.join(creation_info['error_messages'])}" + logger.error(error_msg) + + # Clean up metadata on failure + await self.table_store.delete_collection(user, collection) + + # Clean up tracking + del self.pending_deletions[creation_key] + + raise RuntimeError(error_msg) + + # Clean up tracking + del self.pending_deletions[creation_key] + logger.info(f"Collection {creation_key} auto-created successfully in all storage backends") + except Exception as e: logger.error(f"Error ensuring collection exists: {e}") - # Don't fail the operation if collection creation fails - # This maintains backward compatibility + raise e async def list_collections(self, request: CollectionManagementRequest) -> CollectionManagementResponse: """ @@ -154,6 +208,67 @@ class CollectionManager: tags=tags ) + # Broadcast collection creation to all storage backends + creation_key = (request.user, request.collection) + logger.info(f"Broadcasting create-collection for {creation_key}") + + self.pending_deletions[creation_key] = { + "responses_pending": 3, # vector, object, triples + "responses_received": [], + "all_successful": True, + "error_messages": [], + "deletion_complete": asyncio.Event() + } + + storage_request = StorageManagementRequest( + operation="create-collection", + user=request.user, + collection=request.collection + ) + + # Send creation requests to all storage types + if self.vector_storage_producer: + await self.vector_storage_producer.send(storage_request) + if self.object_storage_producer: + await self.object_storage_producer.send(storage_request) + if self.triples_storage_producer: + await self.triples_storage_producer.send(storage_request) + + # Wait for all storage creations to complete (with timeout) + creation_info = self.pending_deletions[creation_key] + try: + await asyncio.wait_for( + creation_info["deletion_complete"].wait(), + timeout=30.0 # 30 second timeout + ) + except asyncio.TimeoutError: + logger.error(f"Timeout waiting for storage creation responses for {creation_key}") + creation_info["all_successful"] = False + creation_info["error_messages"].append("Timeout waiting for storage creation") + + # Check if all creations succeeded + if not creation_info["all_successful"]: + error_msg = f"Storage creation failed: {'; '.join(creation_info['error_messages'])}" + logger.error(error_msg) + + # Clean up metadata on failure + await self.table_store.delete_collection(request.user, request.collection) + + # Clean up tracking + del self.pending_deletions[creation_key] + + return CollectionManagementResponse( + error=Error( + type="storage_creation_error", + message=error_msg + ), + timestamp=datetime.now().isoformat() + ) + + # Clean up tracking + del self.pending_deletions[creation_key] + logger.info(f"Collection {creation_key} created successfully in all storage backends") + # Get the newly created collection for response created_collection = await self.table_store.get_collection(request.user, request.collection) diff --git a/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py b/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py index b148870c..ad632e2a 100755 --- a/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py @@ -129,7 +129,9 @@ class Processor(TriplesStoreService): logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}") try: - if request.operation == "delete-collection": + if request.operation == "create-collection": + await self.handle_create_collection(request) + elif request.operation == "delete-collection": await self.handle_delete_collection(request) else: response = StorageManagementResponse( @@ -150,6 +152,55 @@ class Processor(TriplesStoreService): ) await self.storage_response_producer.send(response) + async def handle_create_collection(self, request): + """Create a collection in Cassandra triple store""" + try: + # Create or reuse connection for this user's keyspace + if self.table is None or self.table != request.user: + self.tg = None + + try: + if self.cassandra_username and self.cassandra_password: + self.tg = KnowledgeGraph( + hosts=self.cassandra_host, + keyspace=request.user, + username=self.cassandra_username, + password=self.cassandra_password + ) + else: + self.tg = KnowledgeGraph( + hosts=self.cassandra_host, + keyspace=request.user, + ) + except Exception as e: + logger.error(f"Failed to connect to Cassandra for user {request.user}: {e}") + raise + + self.table = request.user + + # Create collection using the built-in method + logger.info(f"Creating collection {request.collection} for user {request.user}") + + if self.tg.collection_exists(request.collection): + logger.info(f"Collection {request.collection} already exists") + else: + self.tg.create_collection(request.collection) + logger.info(f"Created collection {request.collection}") + + # Send success response + response = StorageManagementResponse(error=None) + await self.storage_response_producer.send(response) + + except Exception as e: + logger.error(f"Failed to create collection: {e}", exc_info=True) + response = StorageManagementResponse( + error=Error( + type="creation_error", + message=str(e) + ) + ) + await self.storage_response_producer.send(response) + async def handle_delete_collection(self, request): """Delete all data for a specific collection from the unified triples table""" try: