From 5301960e8b9a6e422d1fe94786189627d6684a14 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 3 Sep 2025 22:12:52 +0100 Subject: [PATCH] Gateway support for objects-query --- .../trustgraph/messaging/__init__.py | 7 ++ .../messaging/translators/__init__.py | 1 + .../messaging/translators/objects_query.py | 79 +++++++++++++++++++ .../trustgraph/gateway/dispatch/manager.py | 2 + .../gateway/dispatch/objects_query.py | 30 +++++++ 5 files changed, 119 insertions(+) create mode 100644 trustgraph-base/trustgraph/messaging/translators/objects_query.py create mode 100644 trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py diff --git a/trustgraph-base/trustgraph/messaging/__init__.py b/trustgraph-base/trustgraph/messaging/__init__.py index 1ed89be7..122acb3a 100644 --- a/trustgraph-base/trustgraph/messaging/__init__.py +++ b/trustgraph-base/trustgraph/messaging/__init__.py @@ -21,6 +21,7 @@ from .translators.embeddings_query import ( DocumentEmbeddingsRequestTranslator, DocumentEmbeddingsResponseTranslator, GraphEmbeddingsRequestTranslator, GraphEmbeddingsResponseTranslator ) +from .translators.objects_query import ObjectsQueryRequestTranslator, ObjectsQueryResponseTranslator # Register all service translators TranslatorRegistry.register_service( @@ -107,6 +108,12 @@ TranslatorRegistry.register_service( GraphEmbeddingsResponseTranslator() ) +TranslatorRegistry.register_service( + "objects-query", + ObjectsQueryRequestTranslator(), + ObjectsQueryResponseTranslator() +) + # 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/__init__.py b/trustgraph-base/trustgraph/messaging/translators/__init__.py index 402b092c..1bed3020 100644 --- a/trustgraph-base/trustgraph/messaging/translators/__init__.py +++ b/trustgraph-base/trustgraph/messaging/translators/__init__.py @@ -17,3 +17,4 @@ from .embeddings_query import ( DocumentEmbeddingsRequestTranslator, DocumentEmbeddingsResponseTranslator, GraphEmbeddingsRequestTranslator, GraphEmbeddingsResponseTranslator ) +from .objects_query import ObjectsQueryRequestTranslator, ObjectsQueryResponseTranslator diff --git a/trustgraph-base/trustgraph/messaging/translators/objects_query.py b/trustgraph-base/trustgraph/messaging/translators/objects_query.py new file mode 100644 index 00000000..a746e0c7 --- /dev/null +++ b/trustgraph-base/trustgraph/messaging/translators/objects_query.py @@ -0,0 +1,79 @@ +from typing import Dict, Any, Tuple, Optional +from ...schema import ObjectsQueryRequest, ObjectsQueryResponse +from .base import MessageTranslator +import json + + +class ObjectsQueryRequestTranslator(MessageTranslator): + """Translator for ObjectsQueryRequest schema objects""" + + def to_pulsar(self, data: Dict[str, Any]) -> ObjectsQueryRequest: + return ObjectsQueryRequest( + user=data.get("user", "trustgraph"), + collection=data.get("collection", "default"), + query=data.get("query", ""), + variables=data.get("variables", {}), + operation_name=data.get("operation_name", None) + ) + + def from_pulsar(self, obj: ObjectsQueryRequest) -> Dict[str, Any]: + result = { + "user": obj.user, + "collection": obj.collection, + "query": obj.query, + "variables": dict(obj.variables) if obj.variables else {} + } + + if obj.operation_name: + result["operation_name"] = obj.operation_name + + return result + + +class ObjectsQueryResponseTranslator(MessageTranslator): + """Translator for ObjectsQueryResponse schema objects""" + + def to_pulsar(self, data: Dict[str, Any]) -> ObjectsQueryResponse: + raise NotImplementedError("Response translation to Pulsar not typically needed") + + def from_pulsar(self, obj: ObjectsQueryResponse) -> Dict[str, Any]: + result = {} + + # Handle GraphQL response data + if obj.data: + try: + result["data"] = json.loads(obj.data) + except json.JSONDecodeError: + result["data"] = obj.data + else: + result["data"] = None + + # Handle GraphQL errors + if obj.errors: + result["errors"] = [] + for error in obj.errors: + error_dict = { + "message": error.message + } + if error.path: + error_dict["path"] = list(error.path) + if error.extensions: + error_dict["extensions"] = dict(error.extensions) + result["errors"].append(error_dict) + + # Handle extensions + if obj.extensions: + result["extensions"] = dict(obj.extensions) + + # 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: ObjectsQueryResponse) -> 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 9ec7b0ab..2dbb4302 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py @@ -19,6 +19,7 @@ from . prompt import PromptRequestor from . graph_rag import GraphRagRequestor from . document_rag import DocumentRagRequestor from . triples_query import TriplesQueryRequestor +from . objects_query import ObjectsQueryRequestor from . embeddings import EmbeddingsRequestor from . graph_embeddings_query import GraphEmbeddingsQueryRequestor from . mcp_tool import McpToolRequestor @@ -50,6 +51,7 @@ request_response_dispatchers = { "embeddings": EmbeddingsRequestor, "graph-embeddings": GraphEmbeddingsQueryRequestor, "triples": TriplesQueryRequestor, + "objects": ObjectsQueryRequestor, } global_dispatchers = { diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py new file mode 100644 index 00000000..2f2535a9 --- /dev/null +++ b/trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py @@ -0,0 +1,30 @@ +from ... schema import ObjectsQueryRequest, ObjectsQueryResponse +from ... messaging import TranslatorRegistry + +from . requestor import ServiceRequestor + +class ObjectsQueryRequestor(ServiceRequestor): + def __init__( + self, pulsar_client, request_queue, response_queue, timeout, + consumer, subscriber, + ): + + super(ObjectsQueryRequestor, self).__init__( + pulsar_client=pulsar_client, + request_queue=request_queue, + response_queue=response_queue, + request_schema=ObjectsQueryRequest, + response_schema=ObjectsQueryResponse, + subscription = subscriber, + consumer_name = consumer, + timeout=timeout, + ) + + self.request_translator = TranslatorRegistry.get_request_translator("objects-query") + self.response_translator = TranslatorRegistry.get_response_translator("objects-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