From e76cdc0f12110b338dd43b03eeff5bfd2bebf106 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 20 Nov 2024 18:38:20 +0000 Subject: [PATCH] Triples query --- test-triples-query-api | 34 +++++++++ trustgraph-flow/scripts/api-gateway | 113 ++++++++++++++++++++++++++++ 2 files changed, 147 insertions(+) create mode 100755 test-triples-query-api diff --git a/test-triples-query-api b/test-triples-query-api new file mode 100755 index 00000000..6caf4d87 --- /dev/null +++ b/test-triples-query-api @@ -0,0 +1,34 @@ +#!/usr/bin/env python3 + +import requests +import json +import sys + +url = "http://localhost:8088/api/v1/" + +############################################################################ + +input = { + "p": "http://www.w3.org/2000/01/rdf-schema#label", + "limit": 10 +} + +resp = requests.post( + f"{url}triples-query", + 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"]) + +sys.exit(0) +############################################################################ + diff --git a/trustgraph-flow/scripts/api-gateway b/trustgraph-flow/scripts/api-gateway index eaf201d7..768e88d5 100755 --- a/trustgraph-flow/scripts/api-gateway +++ b/trustgraph-flow/scripts/api-gateway @@ -27,6 +27,10 @@ from trustgraph.schema import GraphRagQuery, GraphRagResponse from trustgraph.schema import graph_rag_request_queue from trustgraph.schema import graph_rag_response_queue +from trustgraph.schema import TriplesQueryRequest, TriplesQueryResponse, Value +from trustgraph.schema import triples_request_queue +from trustgraph.schema import triples_response_queue + logger = logging.getLogger("api") logger.setLevel(logging.INFO) @@ -131,10 +135,22 @@ class Api: JsonSchema(GraphRagResponse) ) + self.triples_query_out = Publisher( + pulsar_host, triples_request_queue, + schema=JsonSchema(TriplesQueryRequest) + ) + + self.triples_query_in = Subscriber( + pulsar_host, triples_response_queue, + "api-gateway", "api-gateway", + JsonSchema(TriplesQueryResponse) + ) + self.app.add_routes([ web.post("/api/v1/text-completion", self.llm), web.post("/api/v1/prompt", self.prompt), web.post("/api/v1/graph-rag", self.graph_rag), + web.post("/api/v1/triples-query", self.triples_query), ]) async def llm(self, request): @@ -274,6 +290,96 @@ class Api: finally: await self.graph_rag_in.unsubscribe(id) + async def triples_query(self, request): + + id = str(uuid.uuid4()) + + try: + + data = await request.json() + + print(data) + + q = await self.triples_query_in.subscribe(id) + + if "s" in data: + if data["s"].startswith("http:") or data["s"].startswith("https:"): + s = Value(value=data["s"], is_uri=True) + else: + s = Value(value=data["s"], is_uri=True) + else: + s = None + + if "p" in data: + if data["p"].startswith("http:") or data["p"].startswith("https:"): + p = Value(value=data["p"], is_uri=True) + else: + p = Value(value=data["p"], is_uri=True) + else: + p = None + + if "o" in data: + if data["o"].startswith("http:") or data["o"].startswith("https:"): + o = Value(value=data["o"], is_uri=True) + else: + o = Value(value=data["o"], is_uri=True) + else: + o = None + + limit = int(data.get("limit", 10000)) + + await self.triples_query_out.send( + id, + TriplesQueryRequest( + s = s, p = p, o = o, + limit = limit, + user = data.get("user", "trustgraph"), + collection = data.get("collection", "default"), + ) + ) + + try: + resp = await asyncio.wait_for(q.get(), TIME_OUT) + except: + raise RuntimeError("Timeout waiting for response") + + if resp.error: + return web.json_response( + { "error": resp.error.message } + ) + + return web.json_response( + { + "response": [ + { + "s": { + "v": t.s.value, + "e": t.s.is_uri, + }, + "p": { + "v": t.p.value, + "e": t.p.is_uri, + }, + "o": { + "v": t.o.value, + "e": t.o.is_uri, + } + } + for t in resp.triples + ] + } + ) + + except Exception as e: + logging.error(f"Exception: {e}") + + return web.json_response( + { "error": str(e) } + ) + + finally: + await self.graph_rag_in.unsubscribe(id) + async def app_factory(self): self.llm_pub_task = asyncio.create_task(self.llm_in.run()) @@ -285,6 +391,13 @@ class Api: self.graph_rag_pub_task = asyncio.create_task(self.graph_rag_in.run()) self.graph_rag_sub_task = asyncio.create_task(self.graph_rag_out.run()) + self.triples_query_pub_task = asyncio.create_task( + self.triples_query_in.run() + ) + self.triples_query_sub_task = asyncio.create_task( + self.triples_query_out.run() + ) + return self.app def run(self):