From 73538e92caadd6b2da17cf9ca2f31205b813b485 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Sun, 4 May 2025 22:00:00 +0100 Subject: [PATCH] Processing complete, untested --- test-api/test-library-add-processing | 50 +++++++++++++++ test-api/test-library-list-processing | 38 ++++++++++++ .../trustgraph/gateway/dispatch/serialize.py | 4 +- .../trustgraph/librarian/librarian.py | 62 ++++++++++++++++++- .../trustgraph/librarian/table_store.py | 46 +++++++++++++- 5 files changed, 195 insertions(+), 5 deletions(-) create mode 100755 test-api/test-library-add-processing create mode 100755 test-api/test-library-list-processing diff --git a/test-api/test-library-add-processing b/test-api/test-library-add-processing new file mode 100755 index 00000000..f1692b0a --- /dev/null +++ b/test-api/test-library-add-processing @@ -0,0 +1,50 @@ +#!/usr/bin/env python3 + +import requests +import json +import sys +import base64 +import time + +url = "http://localhost:8088/api/v1/" + +############################################################################ + +doc_id = "http://trustgraph.ai/doc/9fdee98b-b259-40ac-bcb9-8e82ccedeb04" + +proc_id = "2714fc72-44ab-45f2-94dd-6773fc336535" + +input = { + "operation": "add-processing", + "processing-metadata": { + "id": proc_id, + "document-id": doc_id, + "time": int(time.time()), + "flow": "0000", + "user": "trustgraph", + "collection": "default", + "tags": ["test"], + } +} + +resp = requests.post( + f"{url}librarian", + json=input, +) + +print(resp.text) +resp = resp.json() + +print(resp) + +if "error" in resp: + print(f"Error: {resp['error']}") + sys.exit(1) + +# print(resp["response"]) +print(resp) + +sys.exit(0) + +############################################################################ + diff --git a/test-api/test-library-list-processing b/test-api/test-library-list-processing new file mode 100755 index 00000000..1d31a572 --- /dev/null +++ b/test-api/test-library-list-processing @@ -0,0 +1,38 @@ +#!/usr/bin/env python3 + +import requests +import json +import sys +import base64 +import time + +url = "http://localhost:8088/api/v1/" + +############################################################################ + +input = { + "operation": "list-processing", + "user": "trustgraph", +} + +resp = requests.post( + f"{url}librarian", + json=input, +) + +print(resp.text) +resp = resp.json() + +print(resp) + +if "error" in resp: + print(f"Error: {resp['error']}") + sys.exit(1) + +# print(resp["response"]) +print(resp) + +sys.exit(0) + +############################################################################ + diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/serialize.py b/trustgraph-flow/trustgraph/gateway/dispatch/serialize.py index 6d7ada54..f6b4e6db 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/serialize.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/serialize.py @@ -118,7 +118,7 @@ def serialize_processing_metadata(message): ret["id"] = message.id if message.id: - ret["document_id"] = message.document_id + ret["document-id"] = message.document_id if message.time: ret["time"] = message.time @@ -154,7 +154,7 @@ def to_processing_metadata(x): return ProcessingMetadata( id = x.get("id", None), - document_id = x.get("document_id", None), + document_id = x.get("document-id", None), time = x.get("time", None), flow = x.get("flow", None), user = x.get("user", None), diff --git a/trustgraph-flow/trustgraph/librarian/librarian.py b/trustgraph-flow/trustgraph/librarian/librarian.py index da716ae5..2a4ff8df 100644 --- a/trustgraph-flow/trustgraph/librarian/librarian.py +++ b/trustgraph-flow/trustgraph/librarian/librarian.py @@ -178,7 +178,40 @@ class Librarian: ) async def add_processing(self, request): - raise RuntimeError("Not implemented") + + print("LIST PROCESSING") + + print("Existence test...") + if await self.table_store.processing_exists( + request.processing_metadata.user, + request.processing_metadata.id + ): + raise RuntimeError("Processing already exists") + + doc = await self.table_store.get_document( + request.processing_metadata.user, + request.processing_metadata.document_id + ) + + object_id = await self.table_store.get_document_object_id( + request.user, + request.document_id + ) + + print("Add processing...") + + await self.table_store.add_processing(request.processing_metadata) + + print("Add complete", flush=True) + + return LibrarianResponse( + error = None, + document_metadata = None, + content = None, + document_metadatas = None, + processing_metadatas = None, + ) + # if document.kind == "application/pdf": # await self.load_document(document) @@ -186,7 +219,30 @@ class Librarian: # await self.load_text(document) async def remove_processing(self, request): - raise RuntimeError("Not implemented") + + print("REMOVING...") + + if not await self.table_store.processing_exists( + request.user, + request.document_id, + ): + raise RuntimeError("Processing object does not exist") + + # Remove doc table row + await self.table_store.remove_processing( + request.user, + request.processing_id + ) + + print("Remove complete", flush=True) + + return LibrarianResponse( + error = None, + document_metadata = None, + content = None, + document_metadatas = None, + processing_metadatas = None, + ) async def list_documents(self, request): @@ -204,6 +260,8 @@ class Librarian: async def list_processing(self, request): + print("LIST PROCESSING") + procs = await self.table_store.list_processing(request.user) print(procs) diff --git a/trustgraph-flow/trustgraph/librarian/table_store.py b/trustgraph-flow/trustgraph/librarian/table_store.py index 47f3cac7..8a5bc3a3 100644 --- a/trustgraph-flow/trustgraph/librarian/table_store.py +++ b/trustgraph-flow/trustgraph/librarian/table_store.py @@ -1,3 +1,4 @@ + from .. schema import LibrarianRequest, LibrarianResponse from .. schema import DocumentMetadata, Error, Triple, Value from .. knowledge import hash @@ -7,8 +8,10 @@ from cassandra.cluster import Cluster from cassandra.auth import PlainTextAuthProvider from cassandra.query import BatchStatement from ssl import SSLContext, PROTOCOL_TLSv1_2 + import uuid import time +import asyncio class TableStore: @@ -224,7 +227,7 @@ class TableStore: ( id, document_id, time, flow, user, collection, - tags, + tags ) VALUES (?, ?, ?, ?, ?, ?, ?) """) @@ -605,6 +608,47 @@ class TableStore: print("Delete complete", flush=True) + async def list_processing(self, user): + + print("LIST processing") + + while True: + + try: + + resp = self.cassandra.execute( + self.list_processing_stmt, + (user,) + ) + + print("OK") + break + + except Exception as e: + print("Exception:", type(e)) + print(f"{e}, retry...", flush=True) + await asyncio.sleep(1) + + + lst = [ + ProcessingMetadata( + id = row[0], + document_id = row[1], + time = int(time.mktime(row[2].timetuple())), + flow = row[3], + user = user, + collection = row[4], + tags = row[5], + ) + for row in resp + ] + + print("OK3") + + print(lst) + + return lst + async def add_graph_embeddings(self, m): when = int(time.time() * 1000)