From 0257d2de6bebcb865c1e96dc81a70ae6eb37f1e7 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Tue, 6 May 2025 23:38:21 +0100 Subject: [PATCH] Knowledge core delete --- trustgraph-flow/trustgraph/cores/knowledge.py | 32 ++++++------- .../trustgraph/gateway/dispatch/knowledge.py | 3 +- .../trustgraph/tables/knowledge.py | 46 +++++++++++++++++++ 3 files changed, 61 insertions(+), 20 deletions(-) diff --git a/trustgraph-flow/trustgraph/cores/knowledge.py b/trustgraph-flow/trustgraph/cores/knowledge.py index b1798af0..adf9b429 100644 --- a/trustgraph-flow/trustgraph/cores/knowledge.py +++ b/trustgraph-flow/trustgraph/cores/knowledge.py @@ -18,28 +18,22 @@ class KnowledgeManager: cassandra_host, cassandra_user, cassandra_password, keyspace ) - async def delete_kg_core(self, request): + async def delete_kg_core(self, request, respond): - print("Updating doc...", flush=True) + print("Deleting core...", flush=True) - # You can't update the document ID, user or kind. + await self.table_store.delete_kg_core( + request.user, request.id + ) - if not await self.table_store.document_exists( - request.document_metadata.user, - request.document_metadata.id - ): - raise RuntimeError("Document does not exist") - - await self.table_store.update_document(request.document_metadata) - - print("Update complete", flush=True) - - return LibrarianResponse( - error = None, - document_metadata = None, - content = None, - document_metadatas = None, - processing_metadatas = None, + await respond( + KnowledgeResponse( + error = None, + ids = None, + eos = False, + triples = None, + graph_embeddings = None, + ) ) async def fetch_kg_core(self, request, respond): diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py b/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py index 5e099a68..2e1ae43a 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py @@ -63,5 +63,6 @@ class KnowledgeRequestor(ServiceRequestor): "eos": True }, True - raise RuntimeError("Unexpected case") + # Empty case, return from successful delete. + return {}, True diff --git a/trustgraph-flow/trustgraph/tables/knowledge.py b/trustgraph-flow/trustgraph/tables/knowledge.py index 50fd9dce..3996b5a7 100644 --- a/trustgraph-flow/trustgraph/tables/knowledge.py +++ b/trustgraph-flow/trustgraph/tables/knowledge.py @@ -181,6 +181,16 @@ class KnowledgeTableStore: WHERE user = ? AND document_id = ? """) + self.delete_triples_stmt = self.cassandra.prepare(""" + DELETE FROM triples + WHERE user = ? AND document_id = ? + """) + + self.delete_graph_embeddings_stmt = self.cassandra.prepare(""" + DELETE FROM graph_embeddings + WHERE user = ? AND document_id = ? + """) + async def add_triples(self, m): when = int(time.time() * 1000) @@ -343,6 +353,42 @@ class KnowledgeTableStore: return lst + async def delete_kg_core(self, user, document_id): + + print("Delete kg cores...") + + while True: + + try: + + resp = self.cassandra.execute( + self.delete_triples_stmt, + (user, document_id) + ) + + break + + except Exception as e: + print("Exception:", type(e)) + print(f"{e}, retry...", flush=True) + await asyncio.sleep(1) + + while True: + + try: + + resp = self.cassandra.execute( + self.delete_graph_embeddings_stmt, + (user, document_id) + ) + + break + + except Exception as e: + print("Exception:", type(e)) + print(f"{e}, retry...", flush=True) + await asyncio.sleep(1) + async def get_triples(self, user, document_id, receiver): print("Get triples...")