mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-24 04:31:02 +02:00
Fixing incomplete porting
This commit is contained in:
parent
3f98e1e0ba
commit
fd749466b6
2 changed files with 3 additions and 19 deletions
|
|
@ -37,8 +37,8 @@ class AsyncProcessor:
|
|||
# Create pub/sub backend via factory
|
||||
self.pubsub_backend = get_pubsub(**params)
|
||||
|
||||
# Keep old PulsarClient for backward compatibility with pulsar_host property
|
||||
self.pulsar_client_object = PulsarClient(**params)
|
||||
# Store pulsar_host for backward compatibility
|
||||
self._pulsar_host = params.get("pulsar_host", "pulsar://pulsar:6650")
|
||||
|
||||
# Initialise metrics, records the parameters
|
||||
ProcessorMetrics(processor = self.id).info({
|
||||
|
|
@ -100,7 +100,6 @@ class AsyncProcessor:
|
|||
# functionality
|
||||
def stop(self):
|
||||
self.pubsub_backend.close()
|
||||
self.pulsar_client.close()
|
||||
self.running = False
|
||||
|
||||
# Returns the pub/sub backend (new interface)
|
||||
|
|
@ -109,11 +108,7 @@ class AsyncProcessor:
|
|||
|
||||
# Returns the pulsar host (backward compatibility)
|
||||
@property
|
||||
def pulsar_host(self): return self.pulsar_client_object.pulsar_host
|
||||
|
||||
# Returns the pulsar client (backward compatibility)
|
||||
@property
|
||||
def pulsar_client(self): return self.pulsar_client_object.client
|
||||
def pulsar_host(self): return self._pulsar_host
|
||||
|
||||
# Register a new event handler for configuration change
|
||||
def register_config_handler(self, handler):
|
||||
|
|
|
|||
|
|
@ -54,17 +54,6 @@ class Api:
|
|||
# Create backend using factory
|
||||
self.pubsub_backend = get_pubsub(**config)
|
||||
|
||||
# Keep pulsar_client for backward compatibility with existing gateway code
|
||||
if self.pulsar_api_key:
|
||||
self.pulsar_client = pulsar.Client(
|
||||
self.pulsar_host, listener_name=self.pulsar_listener,
|
||||
authentication=pulsar.AuthenticationToken(self.pulsar_api_key)
|
||||
)
|
||||
else:
|
||||
self.pulsar_client = pulsar.Client(
|
||||
self.pulsar_host, listener_name=self.pulsar_listener,
|
||||
)
|
||||
|
||||
self.prometheus_url = config.get(
|
||||
"prometheus_url", default_prometheus_url,
|
||||
)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue