From 9aaeec9a616d9cd3691be73f3ea526cef0165234 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Tue, 30 Sep 2025 13:19:26 +0100 Subject: [PATCH] Fix up collection tracking in all processors --- .../direct/milvus_doc_embeddings.py | 16 +++++ .../direct/milvus_graph_embeddings.py | 16 +++++ .../storage/doc_embeddings/milvus/write.py | 42 ++++++++++-- .../storage/doc_embeddings/pinecone/write.py | 66 +++++++++++------- .../storage/graph_embeddings/milvus/write.py | 36 +++++++++- .../graph_embeddings/pinecone/write.py | 65 +++++++++++------- .../storage/objects/cassandra/write.py | 26 ++++++- .../storage/triples/falkordb/write.py | 67 +++++++++++++++++- .../storage/triples/memgraph/write.py | 68 ++++++++++++++++++- 9 files changed, 344 insertions(+), 58 deletions(-) diff --git a/trustgraph-flow/trustgraph/direct/milvus_doc_embeddings.py b/trustgraph-flow/trustgraph/direct/milvus_doc_embeddings.py index 24ac6b23..131c114a 100644 --- a/trustgraph-flow/trustgraph/direct/milvus_doc_embeddings.py +++ b/trustgraph-flow/trustgraph/direct/milvus_doc_embeddings.py @@ -49,6 +49,22 @@ class DocVectors: self.next_reload = time.time() + self.reload_time logger.debug(f"Reload at {self.next_reload}") + def collection_exists(self, user, collection): + """Check if collection exists (dimension-independent check)""" + collection_name = make_safe_collection_name(user, collection, self.prefix) + return self.client.has_collection(collection_name) + + def create_collection(self, user, collection, dimension=384): + """Create collection with default dimension""" + collection_name = make_safe_collection_name(user, collection, self.prefix) + + if self.client.has_collection(collection_name): + logger.info(f"Collection {collection_name} already exists") + return + + self.init_collection(dimension, user, collection) + logger.info(f"Created Milvus collection {collection_name} with dimension {dimension}") + def init_collection(self, dimension, user, collection): collection_name = make_safe_collection_name(user, collection, self.prefix) diff --git a/trustgraph-flow/trustgraph/direct/milvus_graph_embeddings.py b/trustgraph-flow/trustgraph/direct/milvus_graph_embeddings.py index 85292a85..7c2cb55b 100644 --- a/trustgraph-flow/trustgraph/direct/milvus_graph_embeddings.py +++ b/trustgraph-flow/trustgraph/direct/milvus_graph_embeddings.py @@ -49,6 +49,22 @@ class EntityVectors: self.next_reload = time.time() + self.reload_time logger.debug(f"Reload at {self.next_reload}") + def collection_exists(self, user, collection): + """Check if collection exists (dimension-independent check)""" + collection_name = make_safe_collection_name(user, collection, self.prefix) + return self.client.has_collection(collection_name) + + def create_collection(self, user, collection, dimension=384): + """Create collection with default dimension""" + collection_name = make_safe_collection_name(user, collection, self.prefix) + + if self.client.has_collection(collection_name): + logger.info(f"Collection {collection_name} already exists") + return + + self.init_collection(dimension, user, collection) + logger.info(f"Created Milvus collection {collection_name} with dimension {dimension}") + def init_collection(self, dimension, user, collection): collection_name = make_safe_collection_name(user, collection, self.prefix) diff --git a/trustgraph-flow/trustgraph/storage/doc_embeddings/milvus/write.py b/trustgraph-flow/trustgraph/storage/doc_embeddings/milvus/write.py index b2fddeb9..fae3a09a 100755 --- a/trustgraph-flow/trustgraph/storage/doc_embeddings/milvus/write.py +++ b/trustgraph-flow/trustgraph/storage/doc_embeddings/milvus/write.py @@ -68,17 +68,26 @@ class Processor(DocumentEmbeddingsStoreService): async def store_document_embeddings(self, message): + # Validate collection exists before accepting writes + if not self.vecstore.collection_exists(message.metadata.user, message.metadata.collection): + error_msg = ( + f"Collection {message.metadata.collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) + for emb in message.chunks: if emb.chunk is None or emb.chunk == b"": continue - + chunk = emb.chunk.decode("utf-8") if chunk == "": continue for vec in emb.vectors: self.vecstore.insert( - vec, chunk, - message.metadata.user, + vec, chunk, + message.metadata.user, message.metadata.collection ) @@ -99,7 +108,9 @@ class Processor(DocumentEmbeddingsStoreService): 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( @@ -120,6 +131,29 @@ class Processor(DocumentEmbeddingsStoreService): ) await self.storage_response_producer.send(response) + async def handle_create_collection(self, request): + """Create a Milvus collection for document embeddings""" + try: + if self.vecstore.collection_exists(request.user, request.collection): + logger.info(f"Collection {request.user}/{request.collection} already exists") + else: + self.vecstore.create_collection(request.user, request.collection) + logger.info(f"Created collection {request.user}/{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 the collection for document embeddings""" try: diff --git a/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py b/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py index b0a318f7..f940ce2a 100644 --- a/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py +++ b/trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py @@ -123,36 +123,28 @@ class Processor(DocumentEmbeddingsStoreService): async def store_document_embeddings(self, message): + index_name = ( + "d-" + message.metadata.user + "-" + message.metadata.collection + ) + + # Validate collection exists before accepting writes + if not self.pinecone.has_index(index_name): + error_msg = ( + f"Collection {message.metadata.collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) + for emb in message.chunks: if emb.chunk is None or emb.chunk == b"": continue - + chunk = emb.chunk.decode("utf-8") if chunk == "": continue for vec in emb.vectors: - dim = len(vec) - index_name = ( - "d-" + message.metadata.user + "-" + message.metadata.collection - ) - - if index_name != self.last_index_name: - - if not self.pinecone.has_index(index_name): - - try: - - self.create_index(index_name, dim) - - except Exception as e: - logger.error("Pinecone index creation failed") - raise e - - logger.info(f"Index {index_name} created") - - self.last_index_name = index_name - index = self.pinecone.Index(index_name) # Generate unique ID for each vector @@ -204,7 +196,9 @@ class Processor(DocumentEmbeddingsStoreService): 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( @@ -225,6 +219,32 @@ class Processor(DocumentEmbeddingsStoreService): ) await self.storage_response_producer.send(response) + async def handle_create_collection(self, request): + """Create a Pinecone index for document embeddings""" + try: + index_name = f"d-{request.user}-{request.collection}" + + if self.pinecone.has_index(index_name): + logger.info(f"Pinecone index {index_name} already exists") + else: + # Create with default dimension - will need to be recreated if dimension doesn't match + self.create_index(index_name, dim=384) + logger.info(f"Created Pinecone index: {index_name}") + + # 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 the collection for document embeddings""" try: diff --git a/trustgraph-flow/trustgraph/storage/graph_embeddings/milvus/write.py b/trustgraph-flow/trustgraph/storage/graph_embeddings/milvus/write.py index 3058b452..7ccd027b 100755 --- a/trustgraph-flow/trustgraph/storage/graph_embeddings/milvus/write.py +++ b/trustgraph-flow/trustgraph/storage/graph_embeddings/milvus/write.py @@ -68,6 +68,15 @@ class Processor(GraphEmbeddingsStoreService): async def store_graph_embeddings(self, message): + # Validate collection exists before accepting writes + if not self.vecstore.collection_exists(message.metadata.user, message.metadata.collection): + error_msg = ( + f"Collection {message.metadata.collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) + for entity in message.entities: if entity.entity.value != "" and entity.entity.value is not None: @@ -95,7 +104,9 @@ class Processor(GraphEmbeddingsStoreService): 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( @@ -116,6 +127,29 @@ class Processor(GraphEmbeddingsStoreService): ) await self.storage_response_producer.send(response) + async def handle_create_collection(self, request): + """Create a Milvus collection for graph embeddings""" + try: + if self.vecstore.collection_exists(request.user, request.collection): + logger.info(f"Collection {request.user}/{request.collection} already exists") + else: + self.vecstore.create_collection(request.user, request.collection) + logger.info(f"Created collection {request.user}/{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 the collection for graph embeddings""" try: diff --git a/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py b/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py index 02d61a89..f97f2a46 100755 --- a/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py +++ b/trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py @@ -123,6 +123,19 @@ class Processor(GraphEmbeddingsStoreService): async def store_graph_embeddings(self, message): + index_name = ( + "t-" + message.metadata.user + "-" + message.metadata.collection + ) + + # Validate collection exists before accepting writes + if not self.pinecone.has_index(index_name): + error_msg = ( + f"Collection {message.metadata.collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) + for entity in message.entities: if entity.entity.value == "" or entity.entity.value is None: @@ -130,28 +143,6 @@ class Processor(GraphEmbeddingsStoreService): for vec in entity.vectors: - dim = len(vec) - - index_name = ( - "t-" + message.metadata.user + "-" + message.metadata.collection - ) - - if index_name != self.last_index_name: - - if not self.pinecone.has_index(index_name): - - try: - - self.create_index(index_name, dim) - - except Exception as e: - logger.error("Pinecone index creation failed") - raise e - - logger.info(f"Index {index_name} created") - - self.last_index_name = index_name - index = self.pinecone.Index(index_name) # Generate unique ID for each vector @@ -203,7 +194,9 @@ class Processor(GraphEmbeddingsStoreService): 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( @@ -224,6 +217,32 @@ class Processor(GraphEmbeddingsStoreService): ) await self.storage_response_producer.send(response) + async def handle_create_collection(self, request): + """Create a Pinecone index for graph embeddings""" + try: + index_name = f"t-{request.user}-{request.collection}" + + if self.pinecone.has_index(index_name): + logger.info(f"Pinecone index {index_name} already exists") + else: + # Create with default dimension - will need to be recreated if dimension doesn't match + self.create_index(index_name, dim=384) + logger.info(f"Created Pinecone index: {index_name}") + + # 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 the collection for graph embeddings""" try: diff --git a/trustgraph-flow/trustgraph/storage/objects/cassandra/write.py b/trustgraph-flow/trustgraph/storage/objects/cassandra/write.py index 809d0261..fdd32a55 100644 --- a/trustgraph-flow/trustgraph/storage/objects/cassandra/write.py +++ b/trustgraph-flow/trustgraph/storage/objects/cassandra/write.py @@ -348,16 +348,36 @@ class Processor(FlowProcessor): async def on_object(self, msg, consumer, flow): """Process incoming ExtractedObject and store in Cassandra""" - + obj = msg.value() logger.info(f"Storing {len(obj.values)} objects for schema {obj.schema_name} from {obj.metadata.id}") - + + # Validate collection/keyspace exists before accepting writes + safe_keyspace = self.sanitize_name(obj.metadata.user) + if safe_keyspace not in self.known_keyspaces: + # Check if keyspace actually exists in Cassandra + self.connect_cassandra() + check_keyspace_cql = """ + SELECT keyspace_name FROM system_schema.keyspaces + WHERE keyspace_name = %s + """ + result = self.session.execute(check_keyspace_cql, (safe_keyspace,)) + if not result.one(): + error_msg = ( + f"Collection {obj.metadata.collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) + # Cache it if it exists + self.known_keyspaces.add(safe_keyspace) + # Get schema definition schema = self.schemas.get(obj.schema_name) if not schema: logger.warning(f"No schema found for {obj.schema_name} - skipping") return - + # Ensure table exists keyspace = obj.metadata.user table_name = obj.schema_name diff --git a/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py b/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py index 8687271a..d0800b67 100755 --- a/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/falkordb/write.py @@ -152,11 +152,43 @@ class Processor(TriplesStoreService): time=res.run_time_ms )) + def collection_exists(self, user, collection): + """Check if collection metadata node exists""" + result = self.io.query( + "MATCH (c:CollectionMetadata {user: $user, collection: $collection}) " + "RETURN c LIMIT 1", + params={"user": user, "collection": collection} + ) + return result.result_set is not None and len(result.result_set) > 0 + + def create_collection(self, user, collection): + """Create collection metadata node""" + import datetime + self.io.query( + "MERGE (c:CollectionMetadata {user: $user, collection: $collection}) " + "SET c.created_at = $created_at", + params={ + "user": user, + "collection": collection, + "created_at": datetime.datetime.now().isoformat() + } + ) + logger.info(f"Created collection metadata node for {user}/{collection}") + async def store_triples(self, message): # Extract user and collection from metadata user = message.metadata.user if message.metadata.user else "default" collection = message.metadata.collection if message.metadata.collection else "default" + # Validate collection exists before accepting writes + if not self.collection_exists(user, collection): + error_msg = ( + f"Collection {collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) + for t in message.triples: self.create_node(t.s.value, user, collection) @@ -197,7 +229,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( @@ -218,6 +252,29 @@ class Processor(TriplesStoreService): ) await self.storage_response_producer.send(response) + async def handle_create_collection(self, request): + """Create collection metadata in FalkorDB""" + try: + if self.collection_exists(request.user, request.collection): + logger.info(f"Collection {request.user}/{request.collection} already exists") + else: + self.create_collection(request.user, request.collection) + logger.info(f"Created collection {request.user}/{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 the collection for FalkorDB triples""" try: @@ -232,7 +289,13 @@ class Processor(TriplesStoreService): params={"user": request.user, "collection": request.collection} ) - logger.info(f"Deleted {node_result.nodes_deleted} nodes and {literal_result.nodes_deleted} literals for collection {request.user}/{request.collection}") + # Delete collection metadata node + metadata_result = self.io.query( + "MATCH (c:CollectionMetadata {user: $user, collection: $collection}) DELETE c", + params={"user": request.user, "collection": request.collection} + ) + + logger.info(f"Deleted {node_result.nodes_deleted} nodes, {literal_result.nodes_deleted} literals, and {metadata_result.nodes_deleted} metadata nodes for collection {request.user}/{request.collection}") # Send success response response = StorageManagementResponse( diff --git a/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py b/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py index 24083166..84248952 100755 --- a/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/memgraph/write.py @@ -267,12 +267,43 @@ class Processor(TriplesStoreService): src=t.s.value, dest=t.o.value, uri=t.p.value, user=user, collection=collection, ) + def collection_exists(self, user, collection): + """Check if collection metadata node exists""" + with self.io.session(database=self.db) as session: + result = session.run( + "MATCH (c:CollectionMetadata {user: $user, collection: $collection}) " + "RETURN c LIMIT 1", + user=user, collection=collection + ) + return bool(list(result)) + + def create_collection(self, user, collection): + """Create collection metadata node""" + import datetime + with self.io.session(database=self.db) as session: + session.run( + "MERGE (c:CollectionMetadata {user: $user, collection: $collection}) " + "SET c.created_at = $created_at", + user=user, collection=collection, + created_at=datetime.datetime.now().isoformat() + ) + logger.info(f"Created collection metadata node for {user}/{collection}") + async def store_triples(self, message): # Extract user and collection from metadata user = message.metadata.user if message.metadata.user else "default" collection = message.metadata.collection if message.metadata.collection else "default" + # Validate collection exists before accepting writes + if not self.collection_exists(user, collection): + error_msg = ( + f"Collection {collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) + for t in message.triples: self.create_node(t.s.value, user, collection) @@ -329,7 +360,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( @@ -350,6 +383,29 @@ class Processor(TriplesStoreService): ) await self.storage_response_producer.send(response) + async def handle_create_collection(self, request): + """Create collection metadata in Memgraph""" + try: + if self.collection_exists(request.user, request.collection): + logger.info(f"Collection {request.user}/{request.collection} already exists") + else: + self.create_collection(request.user, request.collection) + logger.info(f"Created collection {request.user}/{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""" try: @@ -370,9 +426,17 @@ class Processor(TriplesStoreService): ) literals_deleted = literal_result.consume().counters.nodes_deleted + # Delete collection metadata node + metadata_result = session.run( + "MATCH (c:CollectionMetadata {user: $user, collection: $collection}) " + "DELETE c", + user=request.user, collection=request.collection + ) + metadata_deleted = metadata_result.consume().counters.nodes_deleted + # Note: Relationships are automatically deleted with DETACH DELETE - logger.info(f"Deleted {nodes_deleted} nodes and {literals_deleted} literals for {request.user}/{request.collection}") + logger.info(f"Deleted {nodes_deleted} nodes, {literals_deleted} literals, and {metadata_deleted} metadata nodes for {request.user}/{request.collection}") # Send success response response = StorageManagementResponse(