From 05e301ac9f112dc1e1923110271999c564edd076 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Tue, 30 Sep 2025 12:53:10 +0100 Subject: [PATCH] Remove implicit collection creation from Cassandra --- .../storage/doc_embeddings/qdrant/write.py | 32 ++++++++----------- .../storage/graph_embeddings/qdrant/write.py | 24 ++++++-------- .../storage/triples/cassandra/write.py | 9 ++++++ 3 files changed, 32 insertions(+), 33 deletions(-) diff --git a/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py b/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py index 8898b8e9..dfb9980f 100644 --- a/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py +++ b/trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py @@ -79,6 +79,20 @@ class Processor(DocumentEmbeddingsStoreService): async def store_document_embeddings(self, message): + # Validate collection exists before accepting writes + collection = ( + "d_" + message.metadata.user + "_" + + message.metadata.collection + ) + + if not self.qdrant.collection_exists(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: chunk = emb.chunk.decode("utf-8") @@ -86,24 +100,6 @@ class Processor(DocumentEmbeddingsStoreService): for vec in emb.vectors: - dim = len(vec) - collection = ( - "d_" + message.metadata.user + "_" + - message.metadata.collection - ) - - if not self.qdrant.collection_exists(collection): - try: - self.qdrant.create_collection( - collection_name=collection, - vectors_config=VectorParams( - size=dim, distance=Distance.COSINE - ), - ) - except Exception as e: - logger.error("Qdrant collection creation failed") - raise e - self.qdrant.upsert( collection_name=collection, points=[ diff --git a/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py b/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py index 0d1b7674..6446f7f5 100755 --- a/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py +++ b/trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py @@ -69,23 +69,19 @@ class Processor(GraphEmbeddingsStoreService): metrics=storage_response_metrics, ) - def get_collection(self, dim, user, collection): - + def get_collection(self, user, collection): + """Get collection name and validate it exists""" cname = ( "t_" + user + "_" + collection ) if not self.qdrant.collection_exists(cname): - try: - self.qdrant.create_collection( - collection_name=cname, - vectors_config=VectorParams( - size=dim, distance=Distance.COSINE - ), - ) - except Exception as e: - logger.error("Qdrant collection creation failed") - raise e + error_msg = ( + f"Collection {collection} does not exist. " + f"Create it first with tg-set-collection." + ) + logger.error(error_msg) + raise ValueError(error_msg) return cname @@ -105,10 +101,8 @@ class Processor(GraphEmbeddingsStoreService): for vec in entity.vectors: - dim = len(vec) - collection = self.get_collection( - dim, message.metadata.user, message.metadata.collection + message.metadata.user, message.metadata.collection ) self.qdrant.upsert( diff --git a/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py b/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py index ad632e2a..6497f95c 100755 --- a/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py +++ b/trustgraph-flow/trustgraph/storage/triples/cassandra/write.py @@ -109,6 +109,15 @@ class Processor(TriplesStoreService): self.table = user + # Validate collection exists before accepting writes + if not self.tg.collection_exists(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 t in message.triples: self.tg.insert( message.metadata.collection,