From 4138165589c0f77ac5213667b71795599d5f362f Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Thu, 4 Sep 2025 15:15:31 +0100 Subject: [PATCH] Gateway support for nlp-query and structured-query --- .../trustgraph/messaging/__init__.py | 14 +++++ .../messaging/translators/nlp_query.py | 47 ++++++++++++++++ .../messaging/translators/structured_query.py | 56 +++++++++++++++++++ .../trustgraph/gateway/dispatch/manager.py | 4 ++ .../trustgraph/gateway/dispatch/nlp_query.py | 30 ++++++++++ .../gateway/dispatch/structured_query.py | 30 ++++++++++ 6 files changed, 181 insertions(+) create mode 100644 trustgraph-base/trustgraph/messaging/translators/nlp_query.py create mode 100644 trustgraph-base/trustgraph/messaging/translators/structured_query.py create mode 100644 trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py create mode 100644 trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py diff --git a/trustgraph-base/trustgraph/messaging/__init__.py b/trustgraph-base/trustgraph/messaging/__init__.py index 122acb3a..6b1aedd2 100644 --- a/trustgraph-base/trustgraph/messaging/__init__.py +++ b/trustgraph-base/trustgraph/messaging/__init__.py @@ -22,6 +22,8 @@ from .translators.embeddings_query import ( GraphEmbeddingsRequestTranslator, GraphEmbeddingsResponseTranslator ) from .translators.objects_query import ObjectsQueryRequestTranslator, ObjectsQueryResponseTranslator +from .translators.nlp_query import QuestionToStructuredQueryRequestTranslator, QuestionToStructuredQueryResponseTranslator +from .translators.structured_query import StructuredQueryRequestTranslator, StructuredQueryResponseTranslator # Register all service translators TranslatorRegistry.register_service( @@ -114,6 +116,18 @@ TranslatorRegistry.register_service( ObjectsQueryResponseTranslator() ) +TranslatorRegistry.register_service( + "nlp-query", + QuestionToStructuredQueryRequestTranslator(), + QuestionToStructuredQueryResponseTranslator() +) + +TranslatorRegistry.register_service( + "structured-query", + StructuredQueryRequestTranslator(), + StructuredQueryResponseTranslator() +) + # Register single-direction translators for document loading TranslatorRegistry.register_request("document", DocumentTranslator()) TranslatorRegistry.register_request("text-document", TextDocumentTranslator()) diff --git a/trustgraph-base/trustgraph/messaging/translators/nlp_query.py b/trustgraph-base/trustgraph/messaging/translators/nlp_query.py new file mode 100644 index 00000000..2c445579 --- /dev/null +++ b/trustgraph-base/trustgraph/messaging/translators/nlp_query.py @@ -0,0 +1,47 @@ +from typing import Dict, Any, Tuple +from ...schema import QuestionToStructuredQueryRequest, QuestionToStructuredQueryResponse +from .base import MessageTranslator + + +class QuestionToStructuredQueryRequestTranslator(MessageTranslator): + """Translator for QuestionToStructuredQueryRequest schema objects""" + + def to_pulsar(self, data: Dict[str, Any]) -> QuestionToStructuredQueryRequest: + return QuestionToStructuredQueryRequest( + question=data.get("question", ""), + max_results=data.get("max_results", 100) + ) + + def from_pulsar(self, obj: QuestionToStructuredQueryRequest) -> Dict[str, Any]: + return { + "question": obj.question, + "max_results": obj.max_results + } + + +class QuestionToStructuredQueryResponseTranslator(MessageTranslator): + """Translator for QuestionToStructuredQueryResponse schema objects""" + + def to_pulsar(self, data: Dict[str, Any]) -> QuestionToStructuredQueryResponse: + raise NotImplementedError("Response translation to Pulsar not typically needed") + + def from_pulsar(self, obj: QuestionToStructuredQueryResponse) -> Dict[str, Any]: + result = { + "graphql_query": obj.graphql_query, + "variables": dict(obj.variables) if obj.variables else {}, + "detected_schemas": list(obj.detected_schemas) if obj.detected_schemas else [], + "confidence": obj.confidence + } + + # Handle system-level error + if obj.error: + result["error"] = { + "type": obj.error.type, + "message": obj.error.message + } + + return result + + def from_response_with_completion(self, obj: QuestionToStructuredQueryResponse) -> Tuple[Dict[str, Any], bool]: + """Returns (response_dict, is_final)""" + return self.from_pulsar(obj), True \ No newline at end of file diff --git a/trustgraph-base/trustgraph/messaging/translators/structured_query.py b/trustgraph-base/trustgraph/messaging/translators/structured_query.py new file mode 100644 index 00000000..c6a8abc8 --- /dev/null +++ b/trustgraph-base/trustgraph/messaging/translators/structured_query.py @@ -0,0 +1,56 @@ +from typing import Dict, Any, Tuple +from ...schema import StructuredQueryRequest, StructuredQueryResponse +from .base import MessageTranslator +import json + + +class StructuredQueryRequestTranslator(MessageTranslator): + """Translator for StructuredQueryRequest schema objects""" + + def to_pulsar(self, data: Dict[str, Any]) -> StructuredQueryRequest: + return StructuredQueryRequest( + question=data.get("question", "") + ) + + def from_pulsar(self, obj: StructuredQueryRequest) -> Dict[str, Any]: + return { + "question": obj.question + } + + +class StructuredQueryResponseTranslator(MessageTranslator): + """Translator for StructuredQueryResponse schema objects""" + + def to_pulsar(self, data: Dict[str, Any]) -> StructuredQueryResponse: + raise NotImplementedError("Response translation to Pulsar not typically needed") + + def from_pulsar(self, obj: StructuredQueryResponse) -> Dict[str, Any]: + result = {} + + # Handle structured query response data + if obj.data: + try: + result["data"] = json.loads(obj.data) + except json.JSONDecodeError: + result["data"] = obj.data + else: + result["data"] = None + + # Handle errors (array of strings) + if obj.errors: + result["errors"] = list(obj.errors) + else: + result["errors"] = [] + + # Handle system-level error + if obj.error: + result["error"] = { + "type": obj.error.type, + "message": obj.error.message + } + + return result + + def from_response_with_completion(self, obj: StructuredQueryResponse) -> Tuple[Dict[str, Any], bool]: + """Returns (response_dict, is_final)""" + return self.from_pulsar(obj), True \ No newline at end of file diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py index 2dbb4302..e1e3f367 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py @@ -20,6 +20,8 @@ from . graph_rag import GraphRagRequestor from . document_rag import DocumentRagRequestor from . triples_query import TriplesQueryRequestor from . objects_query import ObjectsQueryRequestor +from . nlp_query import NLPQueryRequestor +from . structured_query import StructuredQueryRequestor from . embeddings import EmbeddingsRequestor from . graph_embeddings_query import GraphEmbeddingsQueryRequestor from . mcp_tool import McpToolRequestor @@ -52,6 +54,8 @@ request_response_dispatchers = { "graph-embeddings": GraphEmbeddingsQueryRequestor, "triples": TriplesQueryRequestor, "objects": ObjectsQueryRequestor, + "nlp-query": NLPQueryRequestor, + "structured-query": StructuredQueryRequestor, } global_dispatchers = { diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py new file mode 100644 index 00000000..3cf5684a --- /dev/null +++ b/trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py @@ -0,0 +1,30 @@ +from ... schema import QuestionToStructuredQueryRequest, QuestionToStructuredQueryResponse +from ... messaging import TranslatorRegistry + +from . requestor import ServiceRequestor + +class NLPQueryRequestor(ServiceRequestor): + def __init__( + self, pulsar_client, request_queue, response_queue, timeout, + consumer, subscriber, + ): + + super(NLPQueryRequestor, self).__init__( + pulsar_client=pulsar_client, + request_queue=request_queue, + response_queue=response_queue, + request_schema=QuestionToStructuredQueryRequest, + response_schema=QuestionToStructuredQueryResponse, + subscription = subscriber, + consumer_name = consumer, + timeout=timeout, + ) + + self.request_translator = TranslatorRegistry.get_request_translator("nlp-query") + self.response_translator = TranslatorRegistry.get_response_translator("nlp-query") + + def to_request(self, body): + return self.request_translator.to_pulsar(body) + + def from_response(self, message): + return self.response_translator.from_response_with_completion(message) \ No newline at end of file diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py new file mode 100644 index 00000000..f08ef038 --- /dev/null +++ b/trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py @@ -0,0 +1,30 @@ +from ... schema import StructuredQueryRequest, StructuredQueryResponse +from ... messaging import TranslatorRegistry + +from . requestor import ServiceRequestor + +class StructuredQueryRequestor(ServiceRequestor): + def __init__( + self, pulsar_client, request_queue, response_queue, timeout, + consumer, subscriber, + ): + + super(StructuredQueryRequestor, self).__init__( + pulsar_client=pulsar_client, + request_queue=request_queue, + response_queue=response_queue, + request_schema=StructuredQueryRequest, + response_schema=StructuredQueryResponse, + subscription = subscriber, + consumer_name = consumer, + timeout=timeout, + ) + + self.request_translator = TranslatorRegistry.get_request_translator("structured-query") + self.response_translator = TranslatorRegistry.get_response_translator("structured-query") + + def to_request(self, body): + return self.request_translator.to_pulsar(body) + + def from_response(self, message): + return self.response_translator.from_response_with_completion(message) \ No newline at end of file