From 979d7294d36830172f24212d13820d3a06dbfe6e Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 17 Dec 2025 13:23:31 +0000 Subject: [PATCH] Fixing schema API problems --- .../trustgraph/base/pulsar_backend.py | 2 +- trustgraph-base/trustgraph/clients/base.py | 41 +++++++++---------- .../trustgraph/clients/config_client.py | 6 --- 3 files changed, 20 insertions(+), 29 deletions(-) diff --git a/trustgraph-base/trustgraph/base/pulsar_backend.py b/trustgraph-base/trustgraph/base/pulsar_backend.py index c9130c1e..d6d6d1d1 100644 --- a/trustgraph-base/trustgraph/base/pulsar_backend.py +++ b/trustgraph-base/trustgraph/base/pulsar_backend.py @@ -70,7 +70,7 @@ def dict_to_dataclass(data: dict, cls: type) -> Any: actual_type = field_type # Recursively convert nested dataclasses - if is_dataclass(actual_type): + if is_dataclass(actual_type) and isinstance(value, dict): kwargs[key] = dict_to_dataclass(value, actual_type) elif hasattr(actual_type, '__origin__'): # Handle generic types like list[T] or dict[K, V] diff --git a/trustgraph-base/trustgraph/clients/base.py b/trustgraph-base/trustgraph/clients/base.py index 25eac3b7..3a4da6ec 100644 --- a/trustgraph-base/trustgraph/clients/base.py +++ b/trustgraph-base/trustgraph/clients/base.py @@ -7,6 +7,7 @@ import time from pulsar.schema import JsonSchema from .. exceptions import * +from ..base.pubsub import get_pubsub # Default timeout for a request/response. In seconds. DEFAULT_TIMEOUT=300 @@ -39,30 +40,25 @@ class BaseClient: if subscriber == None: subscriber = str(uuid.uuid4()) - if pulsar_api_key: - auth = pulsar.AuthenticationToken(pulsar_api_key) - self.client = pulsar.Client( - pulsar_host, - logger=pulsar.ConsoleLogger(log_level), - authentication=auth, - listener=listener, - ) - else: - self.client = pulsar.Client( - pulsar_host, - logger=pulsar.ConsoleLogger(log_level), - listener_name=listener, - ) + # Create backend using factory + self.backend = get_pubsub( + pulsar_host=pulsar_host, + pulsar_api_key=pulsar_api_key, + pulsar_listener=listener, + pubsub_backend='pulsar' + ) - self.producer = self.client.create_producer( + self.producer = self.backend.create_producer( topic=input_queue, - schema=JsonSchema(input_schema), + schema=input_schema, chunking_enabled=True, ) - self.consumer = self.client.subscribe( - output_queue, subscriber, - schema=JsonSchema(output_schema), + self.consumer = self.backend.create_consumer( + topic=output_queue, + subscription=subscriber, + schema=output_schema, + consumer_type='shared', ) self.input_schema = input_schema @@ -136,10 +132,11 @@ class BaseClient: if hasattr(self, "consumer"): self.consumer.close() - + if hasattr(self, "producer"): self.producer.flush() self.producer.close() - - self.client.close() + + if hasattr(self, "backend"): + self.backend.close() diff --git a/trustgraph-base/trustgraph/clients/config_client.py b/trustgraph-base/trustgraph/clients/config_client.py index ed8c704a..be2bf5b9 100644 --- a/trustgraph-base/trustgraph/clients/config_client.py +++ b/trustgraph-base/trustgraph/clients/config_client.py @@ -64,7 +64,6 @@ class ConfigClient(BaseClient): def get(self, keys, timeout=300): resp = self.call( - id=id, operation="get", keys=[ ConfigKey( @@ -88,7 +87,6 @@ class ConfigClient(BaseClient): def list(self, type, timeout=300): resp = self.call( - id=id, operation="list", type=type, timeout=timeout @@ -99,7 +97,6 @@ class ConfigClient(BaseClient): def getvalues(self, type, timeout=300): resp = self.call( - id=id, operation="getvalues", type=type, timeout=timeout @@ -117,7 +114,6 @@ class ConfigClient(BaseClient): def delete(self, keys, timeout=300): resp = self.call( - id=id, operation="delete", keys=[ ConfigKey( @@ -134,7 +130,6 @@ class ConfigClient(BaseClient): def put(self, values, timeout=300): resp = self.call( - id=id, operation="put", values=[ ConfigValue( @@ -152,7 +147,6 @@ class ConfigClient(BaseClient): def config(self, timeout=300): resp = self.call( - id=id, operation="config", timeout=timeout )