Fix up collection tracking in all processors

This commit is contained in:
Cyber MacGeddon 2025-09-30 13:19:26 +01:00
parent d38976526e
commit 9aaeec9a61
9 changed files with 344 additions and 58 deletions

View file

@ -49,6 +49,22 @@ class DocVectors:
self.next_reload = time.time() + self.reload_time self.next_reload = time.time() + self.reload_time
logger.debug(f"Reload at {self.next_reload}") 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): def init_collection(self, dimension, user, collection):
collection_name = make_safe_collection_name(user, collection, self.prefix) collection_name = make_safe_collection_name(user, collection, self.prefix)

View file

@ -49,6 +49,22 @@ class EntityVectors:
self.next_reload = time.time() + self.reload_time self.next_reload = time.time() + self.reload_time
logger.debug(f"Reload at {self.next_reload}") 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): def init_collection(self, dimension, user, collection):
collection_name = make_safe_collection_name(user, collection, self.prefix) collection_name = make_safe_collection_name(user, collection, self.prefix)

View file

@ -68,17 +68,26 @@ class Processor(DocumentEmbeddingsStoreService):
async def store_document_embeddings(self, message): 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: for emb in message.chunks:
if emb.chunk is None or emb.chunk == b"": continue if emb.chunk is None or emb.chunk == b"": continue
chunk = emb.chunk.decode("utf-8") chunk = emb.chunk.decode("utf-8")
if chunk == "": continue if chunk == "": continue
for vec in emb.vectors: for vec in emb.vectors:
self.vecstore.insert( self.vecstore.insert(
vec, chunk, vec, chunk,
message.metadata.user, message.metadata.user,
message.metadata.collection message.metadata.collection
) )
@ -99,7 +108,9 @@ class Processor(DocumentEmbeddingsStoreService):
logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}") logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}")
try: 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) await self.handle_delete_collection(request)
else: else:
response = StorageManagementResponse( response = StorageManagementResponse(
@ -120,6 +131,29 @@ class Processor(DocumentEmbeddingsStoreService):
) )
await self.storage_response_producer.send(response) 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): async def handle_delete_collection(self, request):
"""Delete the collection for document embeddings""" """Delete the collection for document embeddings"""
try: try:

View file

@ -123,36 +123,28 @@ class Processor(DocumentEmbeddingsStoreService):
async def store_document_embeddings(self, message): 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: for emb in message.chunks:
if emb.chunk is None or emb.chunk == b"": continue if emb.chunk is None or emb.chunk == b"": continue
chunk = emb.chunk.decode("utf-8") chunk = emb.chunk.decode("utf-8")
if chunk == "": continue if chunk == "": continue
for vec in emb.vectors: 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) index = self.pinecone.Index(index_name)
# Generate unique ID for each vector # 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}") logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}")
try: 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) await self.handle_delete_collection(request)
else: else:
response = StorageManagementResponse( response = StorageManagementResponse(
@ -225,6 +219,32 @@ class Processor(DocumentEmbeddingsStoreService):
) )
await self.storage_response_producer.send(response) 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): async def handle_delete_collection(self, request):
"""Delete the collection for document embeddings""" """Delete the collection for document embeddings"""
try: try:

View file

@ -68,6 +68,15 @@ class Processor(GraphEmbeddingsStoreService):
async def store_graph_embeddings(self, message): 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: for entity in message.entities:
if entity.entity.value != "" and entity.entity.value is not None: 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}") logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}")
try: 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) await self.handle_delete_collection(request)
else: else:
response = StorageManagementResponse( response = StorageManagementResponse(
@ -116,6 +127,29 @@ class Processor(GraphEmbeddingsStoreService):
) )
await self.storage_response_producer.send(response) 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): async def handle_delete_collection(self, request):
"""Delete the collection for graph embeddings""" """Delete the collection for graph embeddings"""
try: try:

View file

@ -123,6 +123,19 @@ class Processor(GraphEmbeddingsStoreService):
async def store_graph_embeddings(self, message): 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: for entity in message.entities:
if entity.entity.value == "" or entity.entity.value is None: if entity.entity.value == "" or entity.entity.value is None:
@ -130,28 +143,6 @@ class Processor(GraphEmbeddingsStoreService):
for vec in entity.vectors: 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) index = self.pinecone.Index(index_name)
# Generate unique ID for each vector # 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}") logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}")
try: 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) await self.handle_delete_collection(request)
else: else:
response = StorageManagementResponse( response = StorageManagementResponse(
@ -224,6 +217,32 @@ class Processor(GraphEmbeddingsStoreService):
) )
await self.storage_response_producer.send(response) 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): async def handle_delete_collection(self, request):
"""Delete the collection for graph embeddings""" """Delete the collection for graph embeddings"""
try: try:

View file

@ -348,16 +348,36 @@ class Processor(FlowProcessor):
async def on_object(self, msg, consumer, flow): async def on_object(self, msg, consumer, flow):
"""Process incoming ExtractedObject and store in Cassandra""" """Process incoming ExtractedObject and store in Cassandra"""
obj = msg.value() obj = msg.value()
logger.info(f"Storing {len(obj.values)} objects for schema {obj.schema_name} from {obj.metadata.id}") 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 # Get schema definition
schema = self.schemas.get(obj.schema_name) schema = self.schemas.get(obj.schema_name)
if not schema: if not schema:
logger.warning(f"No schema found for {obj.schema_name} - skipping") logger.warning(f"No schema found for {obj.schema_name} - skipping")
return return
# Ensure table exists # Ensure table exists
keyspace = obj.metadata.user keyspace = obj.metadata.user
table_name = obj.schema_name table_name = obj.schema_name

View file

@ -152,11 +152,43 @@ class Processor(TriplesStoreService):
time=res.run_time_ms 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): async def store_triples(self, message):
# Extract user and collection from metadata # Extract user and collection from metadata
user = message.metadata.user if message.metadata.user else "default" user = message.metadata.user if message.metadata.user else "default"
collection = message.metadata.collection if message.metadata.collection 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: for t in message.triples:
self.create_node(t.s.value, user, collection) 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}") logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}")
try: 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) await self.handle_delete_collection(request)
else: else:
response = StorageManagementResponse( response = StorageManagementResponse(
@ -218,6 +252,29 @@ class Processor(TriplesStoreService):
) )
await self.storage_response_producer.send(response) 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): async def handle_delete_collection(self, request):
"""Delete the collection for FalkorDB triples""" """Delete the collection for FalkorDB triples"""
try: try:
@ -232,7 +289,13 @@ class Processor(TriplesStoreService):
params={"user": request.user, "collection": request.collection} 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 # Send success response
response = StorageManagementResponse( response = StorageManagementResponse(

View file

@ -267,12 +267,43 @@ class Processor(TriplesStoreService):
src=t.s.value, dest=t.o.value, uri=t.p.value, user=user, collection=collection, 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): async def store_triples(self, message):
# Extract user and collection from metadata # Extract user and collection from metadata
user = message.metadata.user if message.metadata.user else "default" user = message.metadata.user if message.metadata.user else "default"
collection = message.metadata.collection if message.metadata.collection 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: for t in message.triples:
self.create_node(t.s.value, user, collection) 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}") logger.info(f"Storage management request: {request.operation} for {request.user}/{request.collection}")
try: 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) await self.handle_delete_collection(request)
else: else:
response = StorageManagementResponse( response = StorageManagementResponse(
@ -350,6 +383,29 @@ class Processor(TriplesStoreService):
) )
await self.storage_response_producer.send(response) 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): async def handle_delete_collection(self, request):
"""Delete all data for a specific collection""" """Delete all data for a specific collection"""
try: try:
@ -370,9 +426,17 @@ class Processor(TriplesStoreService):
) )
literals_deleted = literal_result.consume().counters.nodes_deleted 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 # 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 # Send success response
response = StorageManagementResponse( response = StorageManagementResponse(