From e56fa5276a838b1d96bd54a31f184e990707117e Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Wed, 23 Aug 2023 18:37:32 +0900 Subject: [PATCH 01/22] Refactor tracing middleware --- .../requirements.txt | 2 +- .../tracing_middleware_config.yaml | 2 +- panini/managers/middleware_manager.py | 23 +- panini/middleware/tracing_middleware.py | 354 +++++++++++------- requirements/defaults.txt | 2 +- setup.cfg | 8 +- setup.py | 12 +- 7 files changed, 254 insertions(+), 149 deletions(-) diff --git a/examples/tracing_middleware_compose/requirements.txt b/examples/tracing_middleware_compose/requirements.txt index c7d4d9b..c221167 100644 --- a/examples/tracing_middleware_compose/requirements.txt +++ b/examples/tracing_middleware_compose/requirements.txt @@ -4,4 +4,4 @@ opentelemetry-sdk==1.15.0 opentelemetry-exporter-otlp-proto-grpc==1.15.0 opentelemetry-exporter-prometheus==1.12.0rc1 python-logging-loki==0.3.1 -PyYAML==5.4 \ No newline at end of file +PyYAML==6.0.1 \ No newline at end of file diff --git a/examples/tracing_middleware_compose/tracing_middleware_config.yaml b/examples/tracing_middleware_compose/tracing_middleware_config.yaml index 6d57e4b..1946e3b 100644 --- a/examples/tracing_middleware_compose/tracing_middleware_config.yaml +++ b/examples/tracing_middleware_compose/tracing_middleware_config.yaml @@ -1,6 +1,6 @@ service_name: abcdef exporter_config: - endpoint: localhost:4317 + endpoint: open_telemetry:4317 insecure: yes headers: null timeout: null diff --git a/panini/managers/middleware_manager.py b/panini/managers/middleware_manager.py index 9677cd3..2d3e785 100644 --- a/panini/managers/middleware_manager.py +++ b/panini/managers/middleware_manager.py @@ -45,6 +45,17 @@ def add_middleware(self, middleware_cls: Type[Middleware], *args, **kwargs): ), f"At least one of the following functions must be implemented: {high_priority_functions + global_functions}" middleware_obj = middleware_cls(*args, **kwargs) + # Check if Tracing Middleware is first in order when adding middlewares + class_name = middleware_obj.__class__.__name__ + middlewares = self._middlewares + + listen_publish_middleware_exists = middlewares.get('listen_publish_middleware') + listen_request_middleware_exists = middlewares.get('listen_request_middleware') + + if (class_name == 'TracingMiddleware' and + listen_publish_middleware_exists and + listen_request_middleware_exists): + raise Exception("TracingMiddleware should be placed first when adding middlewares!") for function_name in high_priority_functions: if function_name in middleware_cls.__dict__: self._middlewares[f"{function_name}_middleware"].append( @@ -86,7 +97,7 @@ def wrap_function_by_middleware(self, function_type: str) -> Callable: def decorator(function: Callable): def wrap_function_by_send_middleware( - func: Callable, single_middleware + func: Callable, single_middleware ) -> Callable: def next_wrapper(subject: str, message, *args, **kwargs): return single_middleware(subject, message, func, *args, **kwargs) @@ -102,7 +113,7 @@ async def async_next_wrapper(subject: str, message, *args, **kwargs): return next_wrapper def wrap_function_by_listen_middleware( - func: Callable, single_middleware + func: Callable, single_middleware ) -> Callable: def next_wrapper(msg): return single_middleware(msg, func) @@ -116,7 +127,7 @@ async def async_next_wrapper(msg): return next_wrapper def build_middleware_wrapper( - func: Callable, middleware_key: str, wrapper: Callable + func: Callable, middleware_key: str, wrapper: Callable ) -> Callable: for middleware in self._middlewares[middleware_key]: func = wrapper(func, middleware) @@ -138,9 +149,9 @@ def build_middleware_wrapper( else: if ( - len(self._middlewares["listen_publish_middleware"]) == 0 - and len(self._middlewares["listen_request_middleware"]) - == 0 + len(self._middlewares["listen_publish_middleware"]) == 0 + and len(self._middlewares["listen_request_middleware"]) + == 0 ): return function diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 3337407..d832f95 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -1,139 +1,217 @@ -import inspect -import uuid -import json -from functools import wraps -from typing import Optional -from nats.aio.msg import Msg -from panini.app import get_app -from opentelemetry import trace -from dataclasses import dataclass -from panini.middleware import Middleware -from panini.managers.event_manager import Listen -from opentelemetry.sdk.trace import TracerProvider, Tracer -from opentelemetry.sdk.trace.export import BatchSpanProcessor -from opentelemetry.sdk.resources import SERVICE_NAME, Resource -from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter -from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator - - -@dataclass -class SpanConfig: - span_name: str - span_attributes: Optional[dict] - - -def register_trace(**decorator_kwargs): # the decorator - def wrapper(f): # a wrapper for the function - if inspect.iscoroutinefunction(f): - @wraps(f) - async def decorated_function(*args, **kwargs): # the decorated function - ctx = trace.get_current_span().get_span_context() - link = trace.Link(ctx) - _tracer = trace.get_tracer(__name__) - with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), - links=[link]): - result = await f(*args, **kwargs) - return result - else: - @wraps(f) - def decorated_function(*args, **kwargs): # the decorated function +try: + from opentelemetry import trace + from opentelemetry.sdk.trace import TracerProvider, Tracer + from opentelemetry.sdk.trace.export import BatchSpanProcessor + from opentelemetry.sdk.resources import SERVICE_NAME, Resource + from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter + from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator + import inspect + import uuid + import json + from functools import wraps + from typing import Optional + from nats.aio.msg import Msg + from panini.app import get_app + from dataclasses import dataclass + from panini.middleware import Middleware + from panini.managers.event_manager import Listen + + + @dataclass + class SpanConfig: + """ + Represents a configuration for a span. + Attributes: + span_name (str): The name of the span. + span_attributes (Optional[dict]): The optional dictionary of attributes for the span. + """ + span_name: str + span_attributes: Optional[dict] + + + def register_trace(**decorator_kwargs): # the decorator + def wrapper(f): # a wrapper for the function + if inspect.iscoroutinefunction(f): + @wraps(f) + async def decorated_function(*args, **kwargs): # the decorated function + ctx = trace.get_current_span().get_span_context() + link = trace.Link(ctx) + _tracer = trace.get_tracer(__name__) + with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), + links=[link]): + result = await f(*args, **kwargs) + return result + else: + @wraps(f) + def decorated_function(*args, **kwargs): # the decorated function + ctx = trace.get_current_span().get_span_context() + link = trace.Link(ctx) + _tracer = trace.get_tracer(__name__) + with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), + links=[link]): + result = f(*args, **kwargs) + return result + + return decorated_function + + return wrapper + + + class OTELTracer: + """ + The OTELTracer class is used to create and configure an OpenTelemetry tracer for distributed tracing. + Attributes: + tracing_config (dict): The configuration dictionary containing the tracing settings. + service_name (str): The name of the service being traced. + _config (dict): The internal configuration dictionary. + _exporter_config (dict): The configuration for the tracer exporter. + _provider_config (dict): The configuration for the tracer provider. + _custom_config (dict): Additional custom configuration provided by the user. + tracer: The OpenTelemetry tracer instance. + Methods: + __init__(self, tracing_config: dict, **kwargs) + Initializes a new instance of the OTELTracer class. + Args: + tracing_config (dict): The configuration dictionary containing the tracing settings. + **kwargs: Additional custom configuration parameters. + create_tracer(self) + Creates and configures the OpenTelemetry tracer. + Returns: + The initialized OpenTelemetry tracer instance. + """ + + def __init__(self, tracing_config: dict, **kwargs): + self._config = tracing_config + self.service_name = self._config["service_name"] + self._exporter_config = self._config["exporter_config"] + self._provider_config = self._config["provider_config"] + self._custom_config = self._config["custom_config"] + self._custom_config.update(**kwargs) + self.tracer = self.create_tracer() + + def create_tracer(self): + resource = Resource(attributes={ + SERVICE_NAME: self.service_name + }) + provider = TracerProvider(resource=resource, **self._provider_config) + processor = BatchSpanProcessor(OTLPSpanExporter(**self._exporter_config)) + provider.add_span_processor(processor) + trace.set_tracer_provider(provider) + return trace.get_tracer(__name__) + + + class TracingMiddleware(Middleware): + """ + Class representing a middleware for tracing requests and events in a Panini application. + Args: + tracing_config (dict): A dictionary containing the configuration parameters for tracing. + Attributes: + _otel_tracer (OTELTracer): The OpenTelemetry Tracer instance. + tracer (Tracer): The Tracer instance extracted from _otel_tracer. + parent (TraceContextTextMapPropagator): The Trace Context TextMap Propagator instance. + Methods: + _create_uuid(): Generates a UUID v4 string. + send_any(subject: str, message: Msg, send_func, *args, **kwargs): Sends a message with tracing information. + wildcard_match(match_key: str, subject: str) -> str: Performs a wildcard match between two strings. + listen_any(msg: Msg, callback): Listens for events with tracing information. + """ + + def __init__( + self, + tracing_config: dict, + **kwargs + ): + self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) + self.tracer: Tracer = self._otel_tracer.tracer + self.parent = TraceContextTextMapPropagator() + super().__init__() + + def _create_uuid(self) -> str: + return uuid.uuid4().hex + + async def send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): + carrier = {} + headers = {} + span_config = kwargs.get("span_config") + use_tracing = kwargs.get("use_tracing", True) + if kwargs.get("use_current_span", False): ctx = trace.get_current_span().get_span_context() - link = trace.Link(ctx) - _tracer = trace.get_tracer(__name__) - with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), - links=[link]): - result = f(*args, **kwargs) - return result - - return decorated_function - - return wrapper - - -class OTELTracer: - def __init__(self, tracing_config: dict, **kwargs): - self._config = tracing_config - self.service_name = self._config["service_name"] - self._exporter_config = self._config["exporter_config"] - self._provider_config = self._config["provider_config"] - self._custom_config = self._config["custom_config"] - self._custom_config.update(**kwargs) - self.tracer = self.create_tracer() - - def create_tracer(self): - resource = Resource(attributes={ - SERVICE_NAME: self.service_name - }) - provider = TracerProvider(resource=resource, **self._provider_config) - processor = BatchSpanProcessor(OTLPSpanExporter(**self._exporter_config)) - provider.add_span_processor(processor) - trace.set_tracer_provider(provider) - return trace.get_tracer(__name__) - - -class TracingMiddleware(Middleware): - def __init__( - self, - tracing_config: dict, - **kwargs - ): - self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) - self.tracer: Tracer = self._otel_tracer.tracer - self.parent = TraceContextTextMapPropagator() - super().__init__() - - def _create_uuid(self) -> str: - return uuid.uuid4().hex - - async def send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): - carrier = {} - headers = {} - span_config = kwargs.get("span_config") - use_tracing = kwargs.get("use_tracing", True) - if kwargs.get("use_current_span", False): - ctx = trace.get_current_span().get_span_context() - link = [trace.Link(ctx)] - else: - link = [] - if not isinstance(span_config, SpanConfig): - span_config = SpanConfig( - span_name=self._create_uuid(), - span_attributes={}) - if use_tracing is True and span_config: - with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: - for attr_key, attr_value in span_config.span_attributes.items(): - span.set_attribute(attr_key, attr_value) - self.parent.inject(carrier=carrier) - headers = { - "tracing_span_name": span_config.span_name, - "tracing_span_carrier": json.dumps(carrier) - } - if "use_tracing" in kwargs: - del kwargs['use_tracing'] - response = await send_func(subject, message, headers=headers) - return response - - async def listen_any(self, msg: Msg, callback): - context = {} - app = get_app() - assert app is not None - listen_obj_list = app._event_manager.subscriptions[msg.subject] - for index in range(0, len(listen_obj_list)): - listen_object: Listen = listen_obj_list[index] - use_tracing = listen_object._meta.get("use_tracing", True) - if use_tracing: - if id(callback) == id(listen_object.callback) and callback.__name__ == listen_object.callback.__name__: - headers = msg.headers - if headers: - context = self.parent.extract(carrier=json.loads(msg.headers.get("tracing_span_carrier", "{}"))) - span_config = listen_object._meta.get("span_config") - if not isinstance(span_config, SpanConfig): - span_config = SpanConfig( - span_name=self._create_uuid(), - span_attributes={}) - with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: - for attr_key, attr_val in span_config.span_attributes.items(): - span.set_attribute(attr_key, attr_val) - response = await callback(msg) - return response - return await callback(msg) + link = [trace.Link(ctx)] + else: + link = [] + if not isinstance(span_config, SpanConfig): + span_config = SpanConfig( + span_name=self._create_uuid(), + span_attributes={}) + if use_tracing is True and span_config: + with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: + for attr_key, attr_value in span_config.span_attributes.items(): + span.set_attribute(attr_key, attr_value) + self.parent.inject(carrier=carrier) + headers = { + "tracing_span_name": span_config.span_name, + "tracing_span_carrier": json.dumps(carrier) + } + if "use_tracing" in kwargs: + del kwargs['use_tracing'] + response = await send_func(subject, message, headers=headers) + return response + + @classmethod + def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: + """Perform a wildcard match between the match_key and the subject""" + split_subject = subject.split('.') + split_key = match_key.split('.') + + # if `>` at the end of match_key + if split_key[-1] == '>': + if len(split_subject) < len(split_key) - 1: # -1 because `>` matches remaining parts + return None + # checking parts before `>` match + for k, s in zip(split_key[:-1], split_subject): + if k != "*" and k != s: + return None + # parts before `>` matched in both + return match_key + + # if not `>` at the end + if len(split_subject) != len(split_key): + return None + + if all(k == "*" or k == s for k, s in zip(split_key, split_subject)): + return match_key + + return None + + async def listen_any(self, msg: Msg, callback): + context = {} + app = get_app() + assert app is not None + for subject in app._event_manager.subscriptions.keys(): + matched_subject = self.wildcard_match(subject, msg.subject) + if not matched_subject: + continue + listen_obj_list = app._event_manager.subscriptions[matched_subject] + listen_object: Listen + for listen_object in listen_obj_list: + use_tracing = listen_object._meta.get("use_tracing", True) + if use_tracing: + if id(callback) == id( + listen_object.callback) and callback.__name__ == listen_object.callback.__name__: + headers = msg.headers + if headers: + context = self.parent.extract( + carrier=json.loads(msg.headers.get("tracing_span_carrier", "{}"))) + span_config = listen_object._meta.get("span_config") + if not isinstance(span_config, SpanConfig): + span_config = SpanConfig( + span_name=self._create_uuid(), + span_attributes={}) + with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: + for attr_key, attr_val in span_config.span_attributes.items(): + span.set_attribute(attr_key, attr_val) + response = await callback(msg) + return response + return await callback(msg) +except ImportError: + raise Exception("Tracing dependencies not found, try to run `$ pip install panini[tracing]`") \ No newline at end of file diff --git a/requirements/defaults.txt b/requirements/defaults.txt index 7068f9d..4b15496 100644 --- a/requirements/defaults.txt +++ b/requirements/defaults.txt @@ -1,7 +1,7 @@ aiohttp==3.8.1 aiohttp-cors==0.7.0 async-timeout==4.0.0 -requests==2.24.0 +requests==2.31.0 six==1.16.0 nats-py==2.0.0 nats-python==0.8.0 diff --git a/setup.cfg b/setup.cfg index 224a779..c7d6a8f 100644 --- a/setup.cfg +++ b/setup.cfg @@ -1,2 +1,8 @@ [metadata] -description-file = README.md \ No newline at end of file +description-file = README.md + +[options.extras_require] +tracing = opentelemetry-api==1.15.0 + opentelemetry-sdk==1.15.0 + opentelemetry-exporter-otlp-proto-grpc==1.15.0 + opentelemetry-exporter-prometheus==1.12.0rc1 \ No newline at end of file diff --git a/setup.py b/setup.py index cdcbf36..70e8770 100644 --- a/setup.py +++ b/setup.py @@ -17,6 +17,13 @@ """ +tracing_dependencies = [ + "opentelemetry-api==1.15.0", + "opentelemetry-sdk==1.15.0", + "opentelemetry-exporter-otlp-proto-grpc==1.15.0", + "opentelemetry-exporter-prometheus==1.12.0rc1" +] + setup( name="panini", version="0.8.1", @@ -53,7 +60,7 @@ "async-timeout==4.0.0", "nats-py==2.2.0", "websocket-client>=1.2.3", - "requests>=2.24.0", + "requests>=2.31.0", "six>=1.15.0", "yarl>=1.6.1", "python-json-logger>=2.0.1", @@ -76,4 +83,7 @@ "Source": "https://github.com/lwinterface/panini/", }, zip_safe=False, + extras_required={ + 'tracing': tracing_dependencies + } ) From 5c77c409fb9e2f26a5ff0809ddd5df475eb457a0 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 24 Aug 2023 01:40:31 +0900 Subject: [PATCH 02/22] Remove checking from middleware_manager, add checks to tracing_middleware --- panini/managers/middleware_manager.py | 11 - panini/middleware/tracing_middleware.py | 415 ++++++++++++------------ 2 files changed, 210 insertions(+), 216 deletions(-) diff --git a/panini/managers/middleware_manager.py b/panini/managers/middleware_manager.py index 2d3e785..361d8dd 100644 --- a/panini/managers/middleware_manager.py +++ b/panini/managers/middleware_manager.py @@ -45,17 +45,6 @@ def add_middleware(self, middleware_cls: Type[Middleware], *args, **kwargs): ), f"At least one of the following functions must be implemented: {high_priority_functions + global_functions}" middleware_obj = middleware_cls(*args, **kwargs) - # Check if Tracing Middleware is first in order when adding middlewares - class_name = middleware_obj.__class__.__name__ - middlewares = self._middlewares - - listen_publish_middleware_exists = middlewares.get('listen_publish_middleware') - listen_request_middleware_exists = middlewares.get('listen_request_middleware') - - if (class_name == 'TracingMiddleware' and - listen_publish_middleware_exists and - listen_request_middleware_exists): - raise Exception("TracingMiddleware should be placed first when adding middlewares!") for function_name in high_priority_functions: if function_name in middleware_cls.__dict__: self._middlewares[f"{function_name}_middleware"].append( diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index d832f95..05d1c1b 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -5,213 +5,218 @@ from opentelemetry.sdk.resources import SERVICE_NAME, Resource from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator - import inspect - import uuid - import json - from functools import wraps - from typing import Optional - from nats.aio.msg import Msg - from panini.app import get_app - from dataclasses import dataclass - from panini.middleware import Middleware - from panini.managers.event_manager import Listen - - - @dataclass - class SpanConfig: - """ - Represents a configuration for a span. - Attributes: - span_name (str): The name of the span. - span_attributes (Optional[dict]): The optional dictionary of attributes for the span. - """ - span_name: str - span_attributes: Optional[dict] - - - def register_trace(**decorator_kwargs): # the decorator - def wrapper(f): # a wrapper for the function - if inspect.iscoroutinefunction(f): - @wraps(f) - async def decorated_function(*args, **kwargs): # the decorated function - ctx = trace.get_current_span().get_span_context() - link = trace.Link(ctx) - _tracer = trace.get_tracer(__name__) - with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), - links=[link]): - result = await f(*args, **kwargs) - return result - else: - @wraps(f) - def decorated_function(*args, **kwargs): # the decorated function - ctx = trace.get_current_span().get_span_context() - link = trace.Link(ctx) - _tracer = trace.get_tracer(__name__) - with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), - links=[link]): - result = f(*args, **kwargs) - return result - - return decorated_function - - return wrapper - - - class OTELTracer: - """ - The OTELTracer class is used to create and configure an OpenTelemetry tracer for distributed tracing. - Attributes: - tracing_config (dict): The configuration dictionary containing the tracing settings. - service_name (str): The name of the service being traced. - _config (dict): The internal configuration dictionary. - _exporter_config (dict): The configuration for the tracer exporter. - _provider_config (dict): The configuration for the tracer provider. - _custom_config (dict): Additional custom configuration provided by the user. - tracer: The OpenTelemetry tracer instance. - Methods: - __init__(self, tracing_config: dict, **kwargs) - Initializes a new instance of the OTELTracer class. - Args: - tracing_config (dict): The configuration dictionary containing the tracing settings. - **kwargs: Additional custom configuration parameters. - create_tracer(self) - Creates and configures the OpenTelemetry tracer. - Returns: - The initialized OpenTelemetry tracer instance. - """ - - def __init__(self, tracing_config: dict, **kwargs): - self._config = tracing_config - self.service_name = self._config["service_name"] - self._exporter_config = self._config["exporter_config"] - self._provider_config = self._config["provider_config"] - self._custom_config = self._config["custom_config"] - self._custom_config.update(**kwargs) - self.tracer = self.create_tracer() - - def create_tracer(self): - resource = Resource(attributes={ - SERVICE_NAME: self.service_name - }) - provider = TracerProvider(resource=resource, **self._provider_config) - processor = BatchSpanProcessor(OTLPSpanExporter(**self._exporter_config)) - provider.add_span_processor(processor) - trace.set_tracer_provider(provider) - return trace.get_tracer(__name__) - - - class TracingMiddleware(Middleware): - """ - Class representing a middleware for tracing requests and events in a Panini application. - Args: - tracing_config (dict): A dictionary containing the configuration parameters for tracing. - Attributes: - _otel_tracer (OTELTracer): The OpenTelemetry Tracer instance. - tracer (Tracer): The Tracer instance extracted from _otel_tracer. - parent (TraceContextTextMapPropagator): The Trace Context TextMap Propagator instance. - Methods: - _create_uuid(): Generates a UUID v4 string. - send_any(subject: str, message: Msg, send_func, *args, **kwargs): Sends a message with tracing information. - wildcard_match(match_key: str, subject: str) -> str: Performs a wildcard match between two strings. - listen_any(msg: Msg, callback): Listens for events with tracing information. - """ - - def __init__( - self, - tracing_config: dict, - **kwargs - ): - self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) - self.tracer: Tracer = self._otel_tracer.tracer - self.parent = TraceContextTextMapPropagator() - super().__init__() - - def _create_uuid(self) -> str: - return uuid.uuid4().hex - - async def send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): - carrier = {} - headers = {} - span_config = kwargs.get("span_config") - use_tracing = kwargs.get("use_tracing", True) - if kwargs.get("use_current_span", False): +except ImportError: + raise Exception("Tracing dependencies not found, try to run `$ pip install panini[tracing]`") +import inspect +import uuid +import json +from functools import wraps +from typing import Optional +from nats.aio.msg import Msg +from panini.app import get_app +from dataclasses import dataclass +from panini.middleware import Middleware +from panini.managers.event_manager import Listen +import warnings + + +@dataclass +class SpanConfig: + """ + Represents a configuration for a span. + Attributes: + span_name (str): The name of the span. + span_attributes (Optional[dict]): The optional dictionary of attributes for the span. + """ + span_name: str + span_attributes: Optional[dict] + + +def register_trace(**decorator_kwargs): # the decorator + def wrapper(f): # a wrapper for the function + if inspect.iscoroutinefunction(f): + @wraps(f) + async def decorated_function(*args, **kwargs): # the decorated function ctx = trace.get_current_span().get_span_context() - link = [trace.Link(ctx)] - else: - link = [] - if not isinstance(span_config, SpanConfig): - span_config = SpanConfig( - span_name=self._create_uuid(), - span_attributes={}) - if use_tracing is True and span_config: - with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: - for attr_key, attr_value in span_config.span_attributes.items(): - span.set_attribute(attr_key, attr_value) - self.parent.inject(carrier=carrier) - headers = { - "tracing_span_name": span_config.span_name, - "tracing_span_carrier": json.dumps(carrier) - } - if "use_tracing" in kwargs: - del kwargs['use_tracing'] - response = await send_func(subject, message, headers=headers) - return response - - @classmethod - def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: - """Perform a wildcard match between the match_key and the subject""" - split_subject = subject.split('.') - split_key = match_key.split('.') - - # if `>` at the end of match_key - if split_key[-1] == '>': - if len(split_subject) < len(split_key) - 1: # -1 because `>` matches remaining parts - return None - # checking parts before `>` match - for k, s in zip(split_key[:-1], split_subject): - if k != "*" and k != s: - return None - # parts before `>` matched in both - return match_key - - # if not `>` at the end - if len(split_subject) != len(split_key): + link = trace.Link(ctx) + _tracer = trace.get_tracer(__name__) + with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), + links=[link]): + result = await f(*args, **kwargs) + return result + else: + @wraps(f) + def decorated_function(*args, **kwargs): # the decorated function + ctx = trace.get_current_span().get_span_context() + link = trace.Link(ctx) + _tracer = trace.get_tracer(__name__) + with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), + links=[link]): + result = f(*args, **kwargs) + return result + + return decorated_function + + return wrapper + + +class OTELTracer: + """ + The OTELTracer class is used to create and configure an OpenTelemetry tracer for distributed tracing. + Attributes: + tracing_config (dict): The configuration dictionary containing the tracing settings. + service_name (str): The name of the service being traced. + _config (dict): The internal configuration dictionary. + _exporter_config (dict): The configuration for the tracer exporter. + _provider_config (dict): The configuration for the tracer provider. + _custom_config (dict): Additional custom configuration provided by the user. + tracer: The OpenTelemetry tracer instance. + Methods: + __init__(self, tracing_config: dict, **kwargs) + Initializes a new instance of the OTELTracer class. + Args: + tracing_config (dict): The configuration dictionary containing the tracing settings. + **kwargs: Additional custom configuration parameters. + create_tracer(self) + Creates and configures the OpenTelemetry tracer. + Returns: + The initialized OpenTelemetry tracer instance. + """ + + def __init__(self, tracing_config: dict, **kwargs): + self._config = tracing_config + self.service_name = self._config["service_name"] + self._exporter_config = self._config["exporter_config"] + self._provider_config = self._config["provider_config"] + self._custom_config = self._config["custom_config"] + self._custom_config.update(**kwargs) + self.tracer = self.create_tracer() + + def create_tracer(self): + resource = Resource(attributes={ + SERVICE_NAME: self.service_name + }) + provider = TracerProvider(resource=resource, **self._provider_config) + processor = BatchSpanProcessor(OTLPSpanExporter(**self._exporter_config)) + provider.add_span_processor(processor) + trace.set_tracer_provider(provider) + return trace.get_tracer(__name__) + + +class TracingMiddleware(Middleware): + """ + Class representing a middleware for tracing requests and events in a Panini application. + Args: + tracing_config (dict): A dictionary containing the configuration parameters for tracing. + Attributes: + _otel_tracer (OTELTracer): The OpenTelemetry Tracer instance. + tracer (Tracer): The Tracer instance extracted from _otel_tracer. + parent (TraceContextTextMapPropagator): The Trace Context TextMap Propagator instance. + Methods: + _create_uuid(): Generates a UUID v4 string. + send_any(subject: str, message: Msg, send_func, *args, **kwargs): Sends a message with tracing information. + wildcard_match(match_key: str, subject: str) -> str: Performs a wildcard match between two strings. + listen_any(msg: Msg, callback): Listens for events with tracing information. + """ + + def __init__( + self, + tracing_config: dict, + **kwargs + ): + self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) + self.tracer: Tracer = self._otel_tracer.tracer + self.parent = TraceContextTextMapPropagator() + super().__init__() + + def _create_uuid(self) -> str: + return uuid.uuid4().hex + + async def send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): + carrier = {} + headers = {} + span_config = kwargs.get("span_config") + use_tracing = kwargs.get("use_tracing", True) + if kwargs.get("use_current_span", False): + ctx = trace.get_current_span().get_span_context() + link = [trace.Link(ctx)] + else: + link = [] + if not isinstance(span_config, SpanConfig): + span_config = SpanConfig( + span_name=self._create_uuid(), + span_attributes={}) + if use_tracing is True and span_config: + with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: + for attr_key, attr_value in span_config.span_attributes.items(): + span.set_attribute(attr_key, attr_value) + self.parent.inject(carrier=carrier) + headers = { + "tracing_span_name": span_config.span_name, + "tracing_span_carrier": json.dumps(carrier) + } + if "use_tracing" in kwargs: + del kwargs['use_tracing'] + response = await send_func(subject, message, headers=headers) + return response + + @classmethod + def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: + """Perform a wildcard match between the match_key and the subject""" + split_subject = subject.split('.') + split_key = match_key.split('.') + + # if `>` at the end of match_key + if split_key[-1] == '>': + if len(split_subject) < len(split_key) - 1: # -1 because `>` matches remaining parts return None + # checking parts before `>` match + for k, s in zip(split_key[:-1], split_subject): + if k != "*" and k != s: + return None + # parts before `>` matched in both + return match_key - if all(k == "*" or k == s for k, s in zip(split_key, split_subject)): - return match_key - + # if not `>` at the end + if len(split_subject) != len(split_key): return None - async def listen_any(self, msg: Msg, callback): - context = {} - app = get_app() - assert app is not None - for subject in app._event_manager.subscriptions.keys(): - matched_subject = self.wildcard_match(subject, msg.subject) - if not matched_subject: - continue - listen_obj_list = app._event_manager.subscriptions[matched_subject] - listen_object: Listen - for listen_object in listen_obj_list: - use_tracing = listen_object._meta.get("use_tracing", True) - if use_tracing: - if id(callback) == id( - listen_object.callback) and callback.__name__ == listen_object.callback.__name__: - headers = msg.headers - if headers: - context = self.parent.extract( - carrier=json.loads(msg.headers.get("tracing_span_carrier", "{}"))) - span_config = listen_object._meta.get("span_config") - if not isinstance(span_config, SpanConfig): - span_config = SpanConfig( - span_name=self._create_uuid(), - span_attributes={}) - with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: - for attr_key, attr_val in span_config.span_attributes.items(): - span.set_attribute(attr_key, attr_val) - response = await callback(msg) - return response - return await callback(msg) -except ImportError: - raise Exception("Tracing dependencies not found, try to run `$ pip install panini[tracing]`") \ No newline at end of file + if all(k == "*" or k == s for k, s in zip(split_key, split_subject)): + return match_key + + return None + + async def listen_any(self, msg: Msg, callback): + context = {} + app = get_app() + assert app is not None + for subject in app._event_manager.subscriptions.keys(): + matched_subject = self.wildcard_match(subject, msg.subject) + if not matched_subject: + continue + listen_obj_list = app._event_manager.subscriptions[matched_subject] + listen_object: Listen + for listen_object in listen_obj_list: + use_tracing = listen_object._meta.get("use_tracing", True) + if use_tracing: + if id(callback) == id( + listen_object.callback) and callback.__name__ == listen_object.callback.__name__: + headers = msg.headers + if headers: + context = self.parent.extract( + carrier=json.loads(msg.headers.get("tracing_span_carrier", "{}"))) + span_config = listen_object._meta.get("span_config") + if not isinstance(span_config, SpanConfig): + span_config = SpanConfig( + span_name=self._create_uuid(), + span_attributes={}) + with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: + for attr_key, attr_val in span_config.span_attributes.items(): + span.set_attribute(attr_key, attr_val) + response = await callback(msg) + return response + else: + warnings.warn( + "TracingMiddleware logic on listener doesn't work, it should be placed first when adding middlewares!") + return await callback(msg) + return await callback(msg) From 7e4f901d81b549783e1786c6d0b110fb7295a674 Mon Sep 17 00:00:00 2001 From: artas728 Date: Thu, 24 Aug 2023 08:51:23 -0700 Subject: [PATCH 03/22] version update --- setup.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/setup.py b/setup.py index 70e8770..2a2b365 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ setup( name="panini", - version="0.8.1", + version="0.8.3", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst", From 506dd9d6d6208849f6ef08081c7713645bdd6f93 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 7 Sep 2023 01:42:06 +0900 Subject: [PATCH 04/22] set attributes for spans --- panini/middleware/tracing_middleware.py | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 05d1c1b..d0f4b07 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -132,7 +132,7 @@ def __init__( def _create_uuid(self) -> str: return uuid.uuid4().hex - async def send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): + async def send_any(self, subject: str, message: dict, send_func, *args, **kwargs): carrier = {} headers = {} span_config = kwargs.get("span_config") @@ -150,6 +150,9 @@ async def send_any(self, subject: str, message: Msg, send_func, *args, **kwargs) with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) + span.set_attribute("nats_subject", subject) + span.set_attribute("nats_message", json.dumps(message)) + span.set_attribute("nats_action", send_func.__name__) self.parent.inject(carrier=carrier) headers = { "tracing_span_name": span_config.span_name, @@ -213,6 +216,9 @@ async def listen_any(self, msg: Msg, callback): with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: for attr_key, attr_val in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_val) + span.set_attribute("nats_subject", subject) + span.set_attribute("nats_message", json.dumps(msg.data)) + span.set_attribute("nats_action", callback.__name__) response = await callback(msg) return response else: From 7105b6d9192ccfc64f8115508a53333d122fed7d Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 7 Sep 2023 01:59:14 +0900 Subject: [PATCH 05/22] add different name for different functions --- panini/middleware/tracing_middleware.py | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index d0f4b07..e5356aa 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -132,7 +132,7 @@ def __init__( def _create_uuid(self) -> str: return uuid.uuid4().hex - async def send_any(self, subject: str, message: dict, send_func, *args, **kwargs): + async def trace_send_any(self, subject: str, message: dict, send_func, *args, **kwargs): carrier = {} headers = {} span_config = kwargs.get("span_config") @@ -163,6 +163,12 @@ async def send_any(self, subject: str, message: dict, send_func, *args, **kwargs response = await send_func(subject, message, headers=headers) return response + async def send_publish(self, subject: str, message, publish_func, *args, **kwargs): + return await self.trace_send_any(subject, message, publish_func, *args, **kwargs) + + async def send_request(self, subject: str, message, request_func, *args, **kwargs): + return await self.trace_send_any(subject, message, request_func, *args, **kwargs) + @classmethod def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: """Perform a wildcard match between the match_key and the subject""" @@ -189,7 +195,7 @@ def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: return None - async def listen_any(self, msg: Msg, callback): + async def trace_listen_any(self, msg: Msg, callback): context = {} app = get_app() assert app is not None @@ -226,3 +232,9 @@ async def listen_any(self, msg: Msg, callback): "TracingMiddleware logic on listener doesn't work, it should be placed first when adding middlewares!") return await callback(msg) return await callback(msg) + + async def listen_publish(self, msg, callback): + return await self.trace_listen_any(msg, callback) + + async def listen_request(self, msg, callback): + return await self.trace_listen_any(msg, callback) From 2d259ae126666c6d6f804c775338a39545e2daa4 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 7 Sep 2023 02:08:12 +0900 Subject: [PATCH 06/22] few fixes for spans --- .../docker-compose.yml | 6 ++++++ panini/middleware/tracing_middleware.py | 19 ++++++++++++++----- 2 files changed, 20 insertions(+), 5 deletions(-) diff --git a/examples/tracing_middleware_compose/docker-compose.yml b/examples/tracing_middleware_compose/docker-compose.yml index 88f8042..cac56aa 100644 --- a/examples/tracing_middleware_compose/docker-compose.yml +++ b/examples/tracing_middleware_compose/docker-compose.yml @@ -26,7 +26,10 @@ services: - main depends_on: - natsserver + - jaeger + - grafana - receiver + - open_telemetry receiver: build: context: . @@ -43,6 +46,9 @@ services: - main depends_on: - natsserver + - jaeger + - grafana + - open_telemetry nats-exporter: image: synadia/prometheus-nats-exporter:0.3.0 command: '-varz "http://nats-server:8222"' diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index e5356aa..daaa181 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -152,7 +152,7 @@ async def trace_send_any(self, subject: str, message: dict, send_func, *args, ** span.set_attribute(attr_key, attr_value) span.set_attribute("nats_subject", subject) span.set_attribute("nats_message", json.dumps(message)) - span.set_attribute("nats_action", send_func.__name__) + span.set_attribute("nats_action", kwargs.get('nats_action', send_func.__name__)) self.parent.inject(carrier=carrier) headers = { "tracing_span_name": span_config.span_name, @@ -164,9 +164,15 @@ async def trace_send_any(self, subject: str, message: dict, send_func, *args, ** return response async def send_publish(self, subject: str, message, publish_func, *args, **kwargs): + kwargs.update({ + "nats_action": "send_publish" + }) return await self.trace_send_any(subject, message, publish_func, *args, **kwargs) async def send_request(self, subject: str, message, request_func, *args, **kwargs): + kwargs.update({ + "nats_action": "send_request" + }) return await self.trace_send_any(subject, message, request_func, *args, **kwargs) @classmethod @@ -195,7 +201,7 @@ def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: return None - async def trace_listen_any(self, msg: Msg, callback): + async def trace_listen_any(self, msg: Msg, callback, nats_action=None): context = {} app = get_app() assert app is not None @@ -224,7 +230,10 @@ async def trace_listen_any(self, msg: Msg, callback): span.set_attribute(attr_key, attr_val) span.set_attribute("nats_subject", subject) span.set_attribute("nats_message", json.dumps(msg.data)) - span.set_attribute("nats_action", callback.__name__) + if nats_action: + span.set_attribute("nats_action", nats_action) + else: + span.set_attribute("nats_action", callback.__name__) response = await callback(msg) return response else: @@ -234,7 +243,7 @@ async def trace_listen_any(self, msg: Msg, callback): return await callback(msg) async def listen_publish(self, msg, callback): - return await self.trace_listen_any(msg, callback) + return await self.trace_listen_any(msg, callback, nats_action="listen_publish") async def listen_request(self, msg, callback): - return await self.trace_listen_any(msg, callback) + return await self.trace_listen_any(msg, callback, nats_action="listen_request") From 6b23ce3a809f7e35f7fe6b576df783098fcf6644 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 7 Sep 2023 16:08:26 +0900 Subject: [PATCH 07/22] Update tracing_middleware.py --- panini/middleware/tracing_middleware.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index daaa181..51104d0 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -132,7 +132,7 @@ def __init__( def _create_uuid(self) -> str: return uuid.uuid4().hex - async def trace_send_any(self, subject: str, message: dict, send_func, *args, **kwargs): + async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): carrier = {} headers = {} span_config = kwargs.get("span_config") From e6eabb7642d1f63a575fbb83160915bfa51fad6e Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 7 Sep 2023 16:10:30 +0900 Subject: [PATCH 08/22] fixes to async logic --- panini/middleware/tracing_middleware.py | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 51104d0..4049539 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -167,13 +167,15 @@ async def send_publish(self, subject: str, message, publish_func, *args, **kwarg kwargs.update({ "nats_action": "send_publish" }) - return await self.trace_send_any(subject, message, publish_func, *args, **kwargs) + response = await self.trace_send_any(subject, message, publish_func, *args, **kwargs) + return response async def send_request(self, subject: str, message, request_func, *args, **kwargs): kwargs.update({ "nats_action": "send_request" }) - return await self.trace_send_any(subject, message, request_func, *args, **kwargs) + response = await self.trace_send_any(subject, message, request_func, *args, **kwargs) + return response @classmethod def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: @@ -243,7 +245,9 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): return await callback(msg) async def listen_publish(self, msg, callback): - return await self.trace_listen_any(msg, callback, nats_action="listen_publish") + response = await self.trace_listen_any(msg, callback, nats_action="listen_publish") + return response async def listen_request(self, msg, callback): - return await self.trace_listen_any(msg, callback, nats_action="listen_request") + response = await self.trace_listen_any(msg, callback, nats_action="listen_request") + return response From 7adabeab4aa6e99ba3b442ed294d40782f00a45b Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 8 Sep 2023 00:05:46 +0900 Subject: [PATCH 09/22] add event instead of tag --- panini/middleware/tracing_middleware.py | 22 +++++++++++++--------- 1 file changed, 13 insertions(+), 9 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 4049539..c0be791 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -150,9 +150,11 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) - span.set_attribute("nats_subject", subject) - span.set_attribute("nats_message", json.dumps(message)) - span.set_attribute("nats_action", kwargs.get('nats_action', send_func.__name__)) + span.add_event("name", {"nast.subject": subject, "nats.message": json.dumps(message), + "nats.action": kwargs.get('nats_action')}) + # span.set_attribute("nats_subject", subject) + # span.set_attribute("nats_message", json.dumps(message)) + # span.set_attribute("nats_action", kwargs.get('nats_action', send_func.__name__)) self.parent.inject(carrier=carrier) headers = { "tracing_span_name": span_config.span_name, @@ -230,12 +232,14 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: for attr_key, attr_val in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_val) - span.set_attribute("nats_subject", subject) - span.set_attribute("nats_message", json.dumps(msg.data)) - if nats_action: - span.set_attribute("nats_action", nats_action) - else: - span.set_attribute("nats_action", callback.__name__) + # span.set_attribute("nats_subject", subject) + # span.set_attribute("nats_message", json.dumps(msg.data)) + # if nats_action: + # span.set_attribute("nats_action", nats_action) + # else: + # span.set_attribute("nats_action", callback.__name__) + span.add_event("name", {"nast.subject": subject, "nats.message": json.dumps(message), + "nats.action": nats_action}) response = await callback(msg) return response else: From 4aa12dd2892846d460716722166dccaf1f5cbded Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 8 Sep 2023 00:07:50 +0900 Subject: [PATCH 10/22] fix typo --- panini/middleware/tracing_middleware.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index c0be791..2e27c08 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -238,7 +238,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): # span.set_attribute("nats_action", nats_action) # else: # span.set_attribute("nats_action", callback.__name__) - span.add_event("name", {"nast.subject": subject, "nats.message": json.dumps(message), + span.add_event("name", {"nast.subject": subject, "nats.message": json.dumps(msg.data), "nats.action": nats_action}) response = await callback(msg) return response From b728d9215b4d9f8b62e6de6ebc7450bb7f65c7e7 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 8 Sep 2023 00:12:23 +0900 Subject: [PATCH 11/22] add nats.subject for tests for cardinality --- panini/middleware/tracing_middleware.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 2e27c08..2cc21d1 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -152,7 +152,7 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k span.set_attribute(attr_key, attr_value) span.add_event("name", {"nast.subject": subject, "nats.message": json.dumps(message), "nats.action": kwargs.get('nats_action')}) - # span.set_attribute("nats_subject", subject) + span.set_attribute("nats.subject", subject) # span.set_attribute("nats_message", json.dumps(message)) # span.set_attribute("nats_action", kwargs.get('nats_action', send_func.__name__)) self.parent.inject(carrier=carrier) @@ -232,7 +232,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: for attr_key, attr_val in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_val) - # span.set_attribute("nats_subject", subject) + span.set_attribute("nats.subject", subject) # span.set_attribute("nats_message", json.dumps(msg.data)) # if nats_action: # span.set_attribute("nats_action", nats_action) From 82bfd7dba32a07bf6678c376fc737662edb3015e Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 15 Sep 2023 17:56:14 +0900 Subject: [PATCH 12/22] Add exception logging in middleware --- panini/middleware/tracing_middleware.py | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 2cc21d1..20f740d 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -1,3 +1,5 @@ +import time + try: from opentelemetry import trace from opentelemetry.sdk.trace import TracerProvider, Tracer @@ -147,10 +149,11 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k span_name=self._create_uuid(), span_attributes={}) if use_tracing is True and span_config: + self.tracer.start_span() with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) - span.add_event("name", {"nast.subject": subject, "nats.message": json.dumps(message), + span.add_event("name", {"nats.subject": subject, "nats.message": json.dumps(message), "nats.action": kwargs.get('nats_action')}) span.set_attribute("nats.subject", subject) # span.set_attribute("nats_message", json.dumps(message)) @@ -162,8 +165,15 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k } if "use_tracing" in kwargs: del kwargs['use_tracing'] - response = await send_func(subject, message, headers=headers) - return response + try: + response = await send_func(subject, message, headers=headers, *args, **kwargs) + return response + except Exception as exc: + if use_tracing is True and span_config: + self.tracer.start_span() + with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: + span.record_exception(exc, timestamp=time.time_ns() // 10_000) + async def send_publish(self, subject: str, message, publish_func, *args, **kwargs): kwargs.update({ From 75c0dfcc038a0e94d650e034acf0a0660cb020c4 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 15 Sep 2023 22:41:26 +0900 Subject: [PATCH 13/22] refactoring tracing middleware 1. Added exception tracing 2. Added ability to add custom events 3. Added request / response data in panini to tracing events for better understanding --- .../tracing_middleware_compose/Dockerfile | 2 +- .../receiver/app/main.py | 24 ++++- .../sender/app/main.py | 20 +++- panini/middleware/tracing_middleware.py | 93 +++++++++++++------ 4 files changed, 108 insertions(+), 31 deletions(-) diff --git a/examples/tracing_middleware_compose/Dockerfile b/examples/tracing_middleware_compose/Dockerfile index be13c2c..6abafad 100644 --- a/examples/tracing_middleware_compose/Dockerfile +++ b/examples/tracing_middleware_compose/Dockerfile @@ -6,7 +6,7 @@ RUN pip install --upgrade pip ADD requirements.txt / -RUN pip install -r requirements.txt +RUN pip install -r requirements.txt --ignore-installed --force-reinstall RUN mkdir /app WORKDIR /app diff --git a/examples/tracing_middleware_compose/receiver/app/main.py b/examples/tracing_middleware_compose/receiver/app/main.py index 9dce23e..5356ae3 100644 --- a/examples/tracing_middleware_compose/receiver/app/main.py +++ b/examples/tracing_middleware_compose/receiver/app/main.py @@ -3,7 +3,7 @@ import yaml from nats.aio.msg import Msg from panini import app as panini_app -from panini.middleware.tracing_middleware import TracingMiddleware, SpanConfig +from panini.middleware.tracing_middleware import TracingMiddleware, SpanConfig, TracingEvent app = panini_app.App( servers=["nats://nats-server:4222"], @@ -18,6 +18,28 @@ async def default_auto_tracing(msg: Msg): return {"result": True} +@app.listen("test.tracing.middleware.with.events") +async def listen_with_events(msg: Msg): + event = TracingEvent( + event_name='tracing_event_01', + event_data={ + "took_ms": 15, + "request_from": "frontend", + "request_type": "POST" + } + ) + return { + "success": True, + "data": [1, 2, 3, 4, 5], + "tracing_events": [event] + } + + +@app.listen("test.tracing.middleware.with.exceptions") +async def listen_with_exception(msg: Msg): + raise Exception("test exception raised!") + + @app.listen("test.tracing.middleware.no_tracing", use_tracing=False) async def restrict_tracing(msg: Msg): return {"result": True} diff --git a/examples/tracing_middleware_compose/sender/app/main.py b/examples/tracing_middleware_compose/sender/app/main.py index 09f782f..f9faaa4 100644 --- a/examples/tracing_middleware_compose/sender/app/main.py +++ b/examples/tracing_middleware_compose/sender/app/main.py @@ -2,7 +2,7 @@ from nats.aio.msg import Msg from panini import app as panini_app -from panini.middleware.tracing_middleware import SpanConfig, TracingMiddleware +from panini.middleware.tracing_middleware import SpanConfig, TracingMiddleware, TracingEvent app = panini_app.App( servers=["nats://nats-server:4222"], @@ -46,12 +46,30 @@ async def custom_span_tracing(): await app.request("test.tracing.middleware.custom_config", {}, span_config=sender_span_config) return {"result": True} + +@app.task() +async def tracing_with_events(): + event = TracingEvent() + await app.publish("test.tracing.middleware.with.events", {"data": [1, 2, 3, 4, 5]}, tracing_events=[event]) + + +@app.task() +async def exception_tracing_in_listen_function(): + await app.publish("test.tracing.middleware.with.exception", {}) + + +@app.task() +async def exception_tracing_in_send_function(): + await app.publish('test', 1) + + ######################## EXAMPLE OF TRACING WITHIN MICROSERVICE ######################## @app.task() async def tracing_within_microservice(): res = await app.request("test.tracing.inside.microservice.query", {"price": 2}) log.info(f"Result from tracing within microservice: {res}") + @app.listen("test.tracing.inside.microservice.*") async def success_or_fail_result(msg: Msg): _, _, _, _, success = msg.subject.split('.') diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 20f740d..8963a3a 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -16,7 +16,7 @@ from typing import Optional from nats.aio.msg import Msg from panini.app import get_app -from dataclasses import dataclass +from dataclasses import dataclass, field from panini.middleware import Middleware from panini.managers.event_manager import Listen import warnings @@ -34,6 +34,30 @@ class SpanConfig: span_attributes: Optional[dict] +@dataclass +class TracingEvent: + """ + TracingEvent class represents an tracing event. + + Attributes: + event_name (str): The name of the event. + event_data (dict): The data associated with the event. + + Usage: + # Create a TracingEvent object + event = TracingEvent("my_event", {"foo": "bar"}) + + # Add the event to list of events + event_list.append(event) + + # Pass event_list to function: + app.request(subject, message, tracing_events=event_list) + + """ + event_name: str + event_data: dict = field(default_factory=dict) + + def register_trace(**decorator_kwargs): # the decorator def wrapper(f): # a wrapper for the function if inspect.iscoroutinefunction(f): @@ -124,6 +148,7 @@ class TracingMiddleware(Middleware): def __init__( self, tracing_config: dict, + service_name: Optional[str] = None, **kwargs ): self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) @@ -139,6 +164,7 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k headers = {} span_config = kwargs.get("span_config") use_tracing = kwargs.get("use_tracing", True) + existing_events = kwargs.pop("tracing_events", []) if kwargs.get("use_current_span", False): ctx = trace.get_current_span().get_span_context() link = [trace.Link(ctx)] @@ -148,32 +174,37 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k span_config = SpanConfig( span_name=self._create_uuid(), span_attributes={}) - if use_tracing is True and span_config: - self.tracer.start_span() + if use_tracing is True: + self.tracer.start_span(name=span_config.span_name) with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) - span.add_event("name", {"nats.subject": subject, "nats.message": json.dumps(message), - "nats.action": kwargs.get('nats_action')}) + existing_event: TracingEvent + for existing_event in existing_events: + span.add_event(existing_event.event_name, existing_event.event_data) + span.add_event("default", {"nats.subject": subject, "nats.message": json.dumps(message), + "nats.action": kwargs.get('nats_action')}) + kwargs.pop('nats_action') span.set_attribute("nats.subject", subject) - # span.set_attribute("nats_message", json.dumps(message)) - # span.set_attribute("nats_action", kwargs.get('nats_action', send_func.__name__)) self.parent.inject(carrier=carrier) headers = { "tracing_span_name": span_config.span_name, "tracing_span_carrier": json.dumps(carrier) } - if "use_tracing" in kwargs: - del kwargs['use_tracing'] - try: - response = await send_func(subject, message, headers=headers, *args, **kwargs) - return response - except Exception as exc: - if use_tracing is True and span_config: - self.tracer.start_span() - with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: - span.record_exception(exc, timestamp=time.time_ns() // 10_000) + kwargs.update({ + 'headers': headers + }) + if "use_tracing" in kwargs: + del kwargs['use_tracing'] + try: + response = await send_func(subject, message, *args, **kwargs) + return response + except Exception as exc: + span.record_exception(exc) + raise exc + response = await send_func(subject, message, *args, **kwargs) + return response async def send_publish(self, subject: str, message, publish_func, *args, **kwargs): kwargs.update({ @@ -243,20 +274,26 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): for attr_key, attr_val in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_val) span.set_attribute("nats.subject", subject) - # span.set_attribute("nats_message", json.dumps(msg.data)) - # if nats_action: - # span.set_attribute("nats_action", nats_action) - # else: - # span.set_attribute("nats_action", callback.__name__) - span.add_event("name", {"nast.subject": subject, "nats.message": json.dumps(msg.data), - "nats.action": nats_action}) - response = await callback(msg) - return response + span.add_event("default", {"nast.subject": subject, "nats.message": json.dumps(msg.data), + "nats.action": nats_action}) + try: + response = await callback(msg) + if 'tracing_events' in response.keys(): + tracing_events = response.pop('tracing_events') + event: TracingEvent + for event in tracing_events: + span.add_event(event.event_name, event.event_data) + span.add_event("listen_response", response) + except Exception as exc: + span.record_exception(exception=exc) + raise exc else: warnings.warn( "TracingMiddleware logic on listener doesn't work, it should be placed first when adding middlewares!") - return await callback(msg) - return await callback(msg) + response = await callback(msg) + else: + response = await callback(msg) + return response async def listen_publish(self, msg, callback): response = await self.trace_listen_any(msg, callback, nats_action="listen_publish") From 2bf0290cf31173abe0a0d8cc64e185035a082897 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 15 Sep 2023 23:05:04 +0900 Subject: [PATCH 14/22] fix orders of events in send_any function --- panini/middleware/tracing_middleware.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 8963a3a..3cb7a42 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -179,11 +179,11 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) + span.add_event("default", {"nats.subject": subject, "nats.message": json.dumps(message), + "nats.action": kwargs.get('nats_action')}) existing_event: TracingEvent for existing_event in existing_events: span.add_event(existing_event.event_name, existing_event.event_data) - span.add_event("default", {"nats.subject": subject, "nats.message": json.dumps(message), - "nats.action": kwargs.get('nats_action')}) kwargs.pop('nats_action') span.set_attribute("nats.subject", subject) self.parent.inject(carrier=carrier) From 653caf37a60566a340daee389c24f0862b5e11b0 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Mon, 18 Sep 2023 15:21:34 +0900 Subject: [PATCH 15/22] fix global type annotation --- panini/middleware/tracing_middleware.py | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 3cb7a42..d2f6808 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -13,7 +13,7 @@ import uuid import json from functools import wraps -from typing import Optional +from typing import Optional, List from nats.aio.msg import Msg from panini.app import get_app from dataclasses import dataclass, field @@ -164,7 +164,7 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k headers = {} span_config = kwargs.get("span_config") use_tracing = kwargs.get("use_tracing", True) - existing_events = kwargs.pop("tracing_events", []) + existing_events: List[TracingEvent] = kwargs.pop("tracing_events", []) if kwargs.get("use_current_span", False): ctx = trace.get_current_span().get_span_context() link = [trace.Link(ctx)] @@ -181,7 +181,6 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k span.set_attribute(attr_key, attr_value) span.add_event("default", {"nats.subject": subject, "nats.message": json.dumps(message), "nats.action": kwargs.get('nats_action')}) - existing_event: TracingEvent for existing_event in existing_events: span.add_event(existing_event.event_name, existing_event.event_data) kwargs.pop('nats_action') @@ -279,8 +278,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): try: response = await callback(msg) if 'tracing_events' in response.keys(): - tracing_events = response.pop('tracing_events') - event: TracingEvent + tracing_events: List[TracingEvent] = response.pop('tracing_events', []) for event in tracing_events: span.add_event(event.event_name, event.event_data) span.add_event("listen_response", response) From 672dcbc25d48e9533e6967eb78846ae716fe1cc9 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 21 Sep 2023 23:21:06 +0900 Subject: [PATCH 16/22] Add ability to include custom middlewares for http server, refactored tracing middleware --- panini/app.py | 8 +- panini/http_server/http_server_app.py | 9 +- panini/middleware/tracing_middleware.py | 161 +++++++++++++++--------- setup.py | 2 +- 4 files changed, 112 insertions(+), 68 deletions(-) diff --git a/panini/app.py b/panini/app.py index 8204d57..90003b2 100644 --- a/panini/app.py +++ b/panini/app.py @@ -135,7 +135,7 @@ def __init__( error = f"App.event_registrar critical error: {str(e)}" raise InitializingEventManagerError(error) - def setup_web_server(self, host=None, port=None, web_app=None, params: dict = None): + def setup_web_server(self, host=None, port=None, web_app=None, params: dict = None, middlewares: list = None): """ Setup server and run with NATS when called app.start() """ @@ -148,7 +148,8 @@ def setup_web_server(self, host=None, port=None, web_app=None, params: dict = No routes=self.http, loop=self.loop, web_app=web_app, - web_server_params=params + web_server_params=params, + middlewares=middlewares ) else: self.http_server = HTTPServer( @@ -156,7 +157,8 @@ def setup_web_server(self, host=None, port=None, web_app=None, params: dict = No loop=self.loop, host=host, port=port, - web_server_params=params + web_server_params=params, + middlewares=middlewares ) def add_filters(self, include: list = None, exclude: list = None): diff --git a/panini/http_server/http_server_app.py b/panini/http_server/http_server_app.py index a5fffcf..ab8ad60 100644 --- a/panini/http_server/http_server_app.py +++ b/panini/http_server/http_server_app.py @@ -14,6 +14,7 @@ def __init__( port: int = None, web_app: web.Application = None, web_server_params=None, + middlewares: list = None, ): if web_server_params is None: web_server_params = {} @@ -27,7 +28,7 @@ def __init__( if web_app: self.web_app = web_app else: - self.web_app = web.Application() + self.web_app = web.Application(middlewares=middlewares) def start_server(self): self._start_server() @@ -35,5 +36,7 @@ def start_server(self): def _start_server(self): self.web_app.add_routes(self.routes) if version.parse(aiohttp.__version__) >= version.parse("3.8.0"): - self.web_server_params['loop'] = self.loop - web.run_app(self.web_app, host=self.host, port=self.port, **self.web_server_params) + self.web_server_params["loop"] = self.loop + web.run_app( + self.web_app, host=self.host, port=self.port, **self.web_server_params + ) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index d2f6808..b84a0ae 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -1,14 +1,16 @@ -import time - try: from opentelemetry import trace from opentelemetry.sdk.trace import TracerProvider, Tracer from opentelemetry.sdk.trace.export import BatchSpanProcessor from opentelemetry.sdk.resources import SERVICE_NAME, Resource from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter - from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator + from opentelemetry.trace.propagation.tracecontext import ( + TraceContextTextMapPropagator, + ) except ImportError: - raise Exception("Tracing dependencies not found, try to run `$ pip install panini[tracing]`") + raise Exception( + "Tracing dependencies not found, try to run `$ pip install panini[tracing]`" + ) import inspect import uuid import json @@ -30,6 +32,7 @@ class SpanConfig: span_name (str): The name of the span. span_attributes (Optional[dict]): The optional dictionary of attributes for the span. """ + span_name: str span_attributes: Optional[dict] @@ -54,6 +57,7 @@ class TracingEvent: app.request(subject, message, tracing_events=event_list) """ + event_name: str event_data: dict = field(default_factory=dict) @@ -61,23 +65,30 @@ class TracingEvent: def register_trace(**decorator_kwargs): # the decorator def wrapper(f): # a wrapper for the function if inspect.iscoroutinefunction(f): + @wraps(f) async def decorated_function(*args, **kwargs): # the decorated function ctx = trace.get_current_span().get_span_context() link = trace.Link(ctx) _tracer = trace.get_tracer(__name__) - with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), - links=[link]): + with _tracer.start_as_current_span( + decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), + links=[link], + ): result = await f(*args, **kwargs) return result + else: + @wraps(f) def decorated_function(*args, **kwargs): # the decorated function ctx = trace.get_current_span().get_span_context() link = trace.Link(ctx) _tracer = trace.get_tracer(__name__) - with _tracer.start_as_current_span(decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), - links=[link]): + with _tracer.start_as_current_span( + decorator_kwargs.get("span_name", f"unknown-{uuid.uuid4().hex}"), + links=[link], + ): result = f(*args, **kwargs) return result @@ -119,9 +130,7 @@ def __init__(self, tracing_config: dict, **kwargs): self.tracer = self.create_tracer() def create_tracer(self): - resource = Resource(attributes={ - SERVICE_NAME: self.service_name - }) + resource = Resource(attributes={SERVICE_NAME: self.service_name}) provider = TracerProvider(resource=resource, **self._provider_config) processor = BatchSpanProcessor(OTLPSpanExporter(**self._exporter_config)) provider.add_span_processor(processor) @@ -145,12 +154,7 @@ class TracingMiddleware(Middleware): listen_any(msg: Msg, callback): Listens for events with tracing information. """ - def __init__( - self, - tracing_config: dict, - service_name: Optional[str] = None, - **kwargs - ): + def __init__(self, tracing_config: dict, **kwargs): self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) self.tracer: Tracer = self._otel_tracer.tracer self.parent = TraceContextTextMapPropagator() @@ -159,11 +163,21 @@ def __init__( def _create_uuid(self) -> str: return uuid.uuid4().hex - async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): + @staticmethod + def extract_context_from_message(msg: Msg): + propagator = TraceContextTextMapPropagator() + context = propagator.extract( + carrier=json.loads(msg.headers.get("tracing_span_carrier", "{}")) + ) + return context + + async def trace_send_any( + self, subject: str, message: Msg, send_func, *args, **kwargs + ): carrier = {} - headers = {} + context = kwargs.pop("tracing_span_carrier", None) span_config = kwargs.get("span_config") - use_tracing = kwargs.get("use_tracing", True) + use_tracing = kwargs.pop("use_tracing", True) existing_events: List[TracingEvent] = kwargs.pop("tracing_events", []) if kwargs.get("use_current_span", False): ctx = trace.get_current_span().get_span_context() @@ -171,63 +185,64 @@ async def trace_send_any(self, subject: str, message: Msg, send_func, *args, **k else: link = [] if not isinstance(span_config, SpanConfig): - span_config = SpanConfig( - span_name=self._create_uuid(), - span_attributes={}) + span_config = SpanConfig(span_name=self._create_uuid(), span_attributes={}) if use_tracing is True: self.tracer.start_span(name=span_config.span_name) - with self.tracer.start_as_current_span(span_config.span_name, links=link) as span: + with self.tracer.start_as_current_span( + span_config.span_name, links=link, context=context + ) as span: for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) - span.add_event("default", {"nats.subject": subject, "nats.message": json.dumps(message), - "nats.action": kwargs.get('nats_action')}) + span.add_event( + kwargs.pop("nats_action"), + { + "nats.subject": subject, + "nats.message": json.dumps(message), + }, + ) for existing_event in existing_events: span.add_event(existing_event.event_name, existing_event.event_data) - kwargs.pop('nats_action') span.set_attribute("nats.subject", subject) self.parent.inject(carrier=carrier) headers = { "tracing_span_name": span_config.span_name, - "tracing_span_carrier": json.dumps(carrier) + "tracing_span_carrier": json.dumps(carrier), } - kwargs.update({ - 'headers': headers - }) - if "use_tracing" in kwargs: - del kwargs['use_tracing'] + kwargs.update({"headers": headers}) try: response = await send_func(subject, message, *args, **kwargs) return response except Exception as exc: span.record_exception(exc) raise exc - response = await send_func(subject, message, *args, **kwargs) return response async def send_publish(self, subject: str, message, publish_func, *args, **kwargs): - kwargs.update({ - "nats_action": "send_publish" - }) - response = await self.trace_send_any(subject, message, publish_func, *args, **kwargs) + kwargs.update({"nats_action": "send_publish"}) + response = await self.trace_send_any( + subject, message, publish_func, *args, **kwargs + ) return response async def send_request(self, subject: str, message, request_func, *args, **kwargs): - kwargs.update({ - "nats_action": "send_request" - }) - response = await self.trace_send_any(subject, message, request_func, *args, **kwargs) + kwargs.update({"nats_action": "send_request"}) + response = await self.trace_send_any( + subject, message, request_func, *args, **kwargs + ) return response @classmethod def wildcard_match(cls, match_key: str, subject: str) -> Optional[str]: """Perform a wildcard match between the match_key and the subject""" - split_subject = subject.split('.') - split_key = match_key.split('.') + split_subject = subject.split(".") + split_key = match_key.split(".") # if `>` at the end of match_key - if split_key[-1] == '>': - if len(split_subject) < len(split_key) - 1: # -1 because `>` matches remaining parts + if split_key[-1] == ">": + if ( + len(split_subject) < len(split_key) - 1 + ): # -1 because `>` matches remaining parts return None # checking parts before `>` match for k, s in zip(split_key[:-1], split_subject): @@ -258,45 +273,69 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): for listen_object in listen_obj_list: use_tracing = listen_object._meta.get("use_tracing", True) if use_tracing: - if id(callback) == id( - listen_object.callback) and callback.__name__ == listen_object.callback.__name__: + if ( + id(callback) == id(listen_object.callback) + and callback.__name__ == listen_object.callback.__name__ + ): headers = msg.headers if headers: context = self.parent.extract( - carrier=json.loads(msg.headers.get("tracing_span_carrier", "{}"))) + carrier=json.loads( + msg.headers.get("tracing_span_carrier", "{}") + ) + ) span_config = listen_object._meta.get("span_config") if not isinstance(span_config, SpanConfig): span_config = SpanConfig( - span_name=self._create_uuid(), - span_attributes={}) - with self.tracer.start_as_current_span(span_config.span_name, context=context) as span: - for attr_key, attr_val in span_config.span_attributes.items(): + span_name=self._create_uuid(), span_attributes={} + ) + with self.tracer.start_as_current_span( + span_config.span_name, context=context + ) as span: + for ( + attr_key, + attr_val, + ) in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_val) span.set_attribute("nats.subject", subject) - span.add_event("default", {"nast.subject": subject, "nats.message": json.dumps(msg.data), - "nats.action": nats_action}) + span.add_event( + nats_action, + { + "nats.subject": subject, + "nats.message": json.dumps(msg.data), + }, + ) try: response = await callback(msg) - if 'tracing_events' in response.keys(): - tracing_events: List[TracingEvent] = response.pop('tracing_events', []) + if "tracing_events" in response.keys(): + tracing_events: List[TracingEvent] = response.pop( + "tracing_events", [] + ) for event in tracing_events: - span.add_event(event.event_name, event.event_data) + span.add_event( + event.event_name, event.event_data + ) span.add_event("listen_response", response) except Exception as exc: span.record_exception(exception=exc) raise exc else: warnings.warn( - "TracingMiddleware logic on listener doesn't work, it should be placed first when adding middlewares!") + "TracingMiddleware logic on listener doesn't work, it should be placed first when adding middlewares!" + ) response = await callback(msg) else: response = await callback(msg) return response async def listen_publish(self, msg, callback): - response = await self.trace_listen_any(msg, callback, nats_action="listen_publish") + response = await self.trace_listen_any( + msg, callback, nats_action="listen_publish" + ) return response async def listen_request(self, msg, callback): - response = await self.trace_listen_any(msg, callback, nats_action="listen_request") + response = await self.trace_listen_any( + msg, callback, nats_action="listen_request" + ) return response diff --git a/setup.py b/setup.py index 2a2b365..38fc6c0 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ setup( name="panini", - version="0.8.3", + version="0.8.3b4", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst", From 8cb94d9755392e209d33f1191ffb6ae6ede56016 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 21 Sep 2023 23:56:16 +0900 Subject: [PATCH 17/22] Bump opentelemetry versions to match django library --- examples/tracing_middleware_compose/requirements.txt | 6 +++--- setup.py | 8 ++++---- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/examples/tracing_middleware_compose/requirements.txt b/examples/tracing_middleware_compose/requirements.txt index c221167..915c693 100644 --- a/examples/tracing_middleware_compose/requirements.txt +++ b/examples/tracing_middleware_compose/requirements.txt @@ -1,7 +1,7 @@ git+https://github.com/i2gor87/panini.git@feature/tracing-middleware -opentelemetry-api==1.15.0 -opentelemetry-sdk==1.15.0 -opentelemetry-exporter-otlp-proto-grpc==1.15.0 +opentelemetry-api==1.19.0 +opentelemetry-sdk==1.19.0 +opentelemetry-exporter-otlp-proto-grpc==1.19.0 opentelemetry-exporter-prometheus==1.12.0rc1 python-logging-loki==0.3.1 PyYAML==6.0.1 \ No newline at end of file diff --git a/setup.py b/setup.py index 38fc6c0..0ce9e1b 100644 --- a/setup.py +++ b/setup.py @@ -18,15 +18,15 @@ """ tracing_dependencies = [ - "opentelemetry-api==1.15.0", - "opentelemetry-sdk==1.15.0", - "opentelemetry-exporter-otlp-proto-grpc==1.15.0", + "opentelemetry-api==1.19.0", + "opentelemetry-sdk==1.19.0", + "opentelemetry-exporter-otlp-proto-grpc==1.19.0", "opentelemetry-exporter-prometheus==1.12.0rc1" ] setup( name="panini", - version="0.8.3b4", + version="0.8.3a5", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst", From 03a600f9f48f6aa728923e65f286a48613ca291a Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 22 Sep 2023 22:03:59 +0900 Subject: [PATCH 18/22] Small fixes and readability Update opentelemetry versions to match django opentelemetry library. Correct span name for readability Added ability to pass common otel_tracer to TracingMiddleware --- panini/middleware/tracing_middleware.py | 17 +++++++++++++---- setup.cfg | 6 +++--- setup.py | 2 +- 3 files changed, 17 insertions(+), 8 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index b84a0ae..ed7b8b5 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -155,14 +155,22 @@ class TracingMiddleware(Middleware): """ def __init__(self, tracing_config: dict, **kwargs): - self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) + otel_tracer = kwargs.get("otel_tracer") + if otel_tracer: + self._otel_tracer = otel_tracer + else: + self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) self.tracer: Tracer = self._otel_tracer.tracer self.parent = TraceContextTextMapPropagator() + self._service_name = tracing_config["service_name"] super().__init__() def _create_uuid(self) -> str: return uuid.uuid4().hex + def _get_span_name(self, action: str, subject:str) -> str: + return f"{self._service_name}: {action.upper()} {subject}" + @staticmethod def extract_context_from_message(msg: Msg): propagator = TraceContextTextMapPropagator() @@ -178,6 +186,7 @@ async def trace_send_any( context = kwargs.pop("tracing_span_carrier", None) span_config = kwargs.get("span_config") use_tracing = kwargs.pop("use_tracing", True) + action = kwargs.pop("nats_action") existing_events: List[TracingEvent] = kwargs.pop("tracing_events", []) if kwargs.get("use_current_span", False): ctx = trace.get_current_span().get_span_context() @@ -185,7 +194,7 @@ async def trace_send_any( else: link = [] if not isinstance(span_config, SpanConfig): - span_config = SpanConfig(span_name=self._create_uuid(), span_attributes={}) + span_config = SpanConfig(span_name=self._get_span_name(action, subject), span_attributes={}) if use_tracing is True: self.tracer.start_span(name=span_config.span_name) with self.tracer.start_as_current_span( @@ -194,7 +203,7 @@ async def trace_send_any( for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) span.add_event( - kwargs.pop("nats_action"), + action, { "nats.subject": subject, "nats.message": json.dumps(message), @@ -287,7 +296,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): span_config = listen_object._meta.get("span_config") if not isinstance(span_config, SpanConfig): span_config = SpanConfig( - span_name=self._create_uuid(), span_attributes={} + span_name=self._get_span_name(nats_action, subject), span_attributes={} ) with self.tracer.start_as_current_span( span_config.span_name, context=context diff --git a/setup.cfg b/setup.cfg index c7d6a8f..1b11281 100644 --- a/setup.cfg +++ b/setup.cfg @@ -2,7 +2,7 @@ description-file = README.md [options.extras_require] -tracing = opentelemetry-api==1.15.0 - opentelemetry-sdk==1.15.0 - opentelemetry-exporter-otlp-proto-grpc==1.15.0 +tracing = opentelemetry-api==1.19.0 + opentelemetry-sdk==1.19.0 + opentelemetry-exporter-otlp-proto-grpc==1.19.0 opentelemetry-exporter-prometheus==1.12.0rc1 \ No newline at end of file diff --git a/setup.py b/setup.py index 0ce9e1b..fd2b9f6 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ setup( name="panini", - version="0.8.3a5", + version="0.8.3a8", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst", From 894f34f14d0a8fbe58db400bf8279bafbb0f57d0 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Fri, 22 Sep 2023 22:54:41 +0900 Subject: [PATCH 19/22] added send response to tracing events --- panini/middleware/tracing_middleware.py | 7 +++++++ setup.py | 2 +- 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index ed7b8b5..e9b960a 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -220,6 +220,13 @@ async def trace_send_any( kwargs.update({"headers": headers}) try: response = await send_func(subject, message, *args, **kwargs) + if response: + span.add_event( + "request_response", + { + "nats.message": response + } + ) return response except Exception as exc: span.record_exception(exc) diff --git a/setup.py b/setup.py index fd2b9f6..fd93d60 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ setup( name="panini", - version="0.8.3a8", + version="0.8.3a9", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst", From 265f2f32b403a4c2659cae48d424817d1eea4be6 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Mon, 25 Sep 2023 20:54:40 +0900 Subject: [PATCH 20/22] fix span name too long --- panini/middleware/tracing_middleware.py | 18 ++++++++---------- setup.py | 2 +- 2 files changed, 9 insertions(+), 11 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index e9b960a..9b61013 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -168,8 +168,8 @@ def __init__(self, tracing_config: dict, **kwargs): def _create_uuid(self) -> str: return uuid.uuid4().hex - def _get_span_name(self, action: str, subject:str) -> str: - return f"{self._service_name}: {action.upper()} {subject}" + def _get_span_name(self, action: str, subject: str) -> str: + return f"{action.upper()} {subject}" @staticmethod def extract_context_from_message(msg: Msg): @@ -194,7 +194,9 @@ async def trace_send_any( else: link = [] if not isinstance(span_config, SpanConfig): - span_config = SpanConfig(span_name=self._get_span_name(action, subject), span_attributes={}) + span_config = SpanConfig( + span_name=self._get_span_name(action, subject), span_attributes={} + ) if use_tracing is True: self.tracer.start_span(name=span_config.span_name) with self.tracer.start_as_current_span( @@ -221,12 +223,7 @@ async def trace_send_any( try: response = await send_func(subject, message, *args, **kwargs) if response: - span.add_event( - "request_response", - { - "nats.message": response - } - ) + span.add_event("request_response", {"nats.message": response}) return response except Exception as exc: span.record_exception(exc) @@ -303,7 +300,8 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): span_config = listen_object._meta.get("span_config") if not isinstance(span_config, SpanConfig): span_config = SpanConfig( - span_name=self._get_span_name(nats_action, subject), span_attributes={} + span_name=self._get_span_name(nats_action, subject), + span_attributes={}, ) with self.tracer.start_as_current_span( span_config.span_name, context=context diff --git a/setup.py b/setup.py index fd93d60..502f657 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ setup( name="panini", - version="0.8.3a9", + version="0.8.3a10", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst", From e0478008527d7eb0417e5868684b51f886a3f49d Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Tue, 26 Sep 2023 00:47:54 +0900 Subject: [PATCH 21/22] added default non-verbose behaviour which could be changed with a flag for selected actions --- panini/middleware/tracing_middleware.py | 11 +++++++---- setup.py | 2 +- 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 9b61013..93d9299 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -182,6 +182,7 @@ def extract_context_from_message(msg: Msg): async def trace_send_any( self, subject: str, message: Msg, send_func, *args, **kwargs ): + verbose = kwargs.pop("verbose", False) carrier = {} context = kwargs.pop("tracing_span_carrier", None) span_config = kwargs.get("span_config") @@ -208,7 +209,7 @@ async def trace_send_any( action, { "nats.subject": subject, - "nats.message": json.dumps(message), + "nats.message": json.dumps(message) if verbose else json.dumps(message)[:300], }, ) for existing_event in existing_events: @@ -222,7 +223,7 @@ async def trace_send_any( kwargs.update({"headers": headers}) try: response = await send_func(subject, message, *args, **kwargs) - if response: + if response and verbose: span.add_event("request_response", {"nats.message": response}) return response except Exception as exc: @@ -286,6 +287,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): for listen_object in listen_obj_list: use_tracing = listen_object._meta.get("use_tracing", True) if use_tracing: + verbose = listen_object._meta.get("verbose", False) if ( id(callback) == id(listen_object.callback) and callback.__name__ == listen_object.callback.__name__ @@ -316,7 +318,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): nats_action, { "nats.subject": subject, - "nats.message": json.dumps(msg.data), + "nats.message": json.dumps(msg.data) if verbose else json.dumps(msg.data)[:300], }, ) try: @@ -329,7 +331,8 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): span.add_event( event.event_name, event.event_data ) - span.add_event("listen_response", response) + if verbose: + span.add_event("listen_response", response) except Exception as exc: span.record_exception(exception=exc) raise exc diff --git a/setup.py b/setup.py index 502f657..d28cd96 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ setup( name="panini", - version="0.8.3a10", + version="0.8.3a11", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst", From e30b69e60f02c82f5aebea5018a752dda136ca51 Mon Sep 17 00:00:00 2001 From: Igor Khvan Date: Thu, 28 Sep 2023 12:42:13 +0500 Subject: [PATCH 22/22] Update PR to latest version --- .../docker-compose.yml | 15 ++++++++-- .../tracing_middleware_compose/jaeger-ui.json | 8 +++++ .../open_telemetry_config.yaml | 29 +++---------------- .../tracing_middleware_compose/prometheus.yml | 7 ++++- .../requirements.txt | 4 ++- panini/middleware/tracing_middleware.py | 11 +++---- setup.py | 2 +- 7 files changed, 40 insertions(+), 36 deletions(-) create mode 100644 examples/tracing_middleware_compose/jaeger-ui.json diff --git a/examples/tracing_middleware_compose/docker-compose.yml b/examples/tracing_middleware_compose/docker-compose.yml index cac56aa..268f95f 100644 --- a/examples/tracing_middleware_compose/docker-compose.yml +++ b/examples/tracing_middleware_compose/docker-compose.yml @@ -29,6 +29,7 @@ services: - jaeger - grafana - receiver + - prometheus - open_telemetry receiver: build: @@ -48,6 +49,7 @@ services: - natsserver - jaeger - grafana + - prometheus - open_telemetry nats-exporter: image: synadia/prometheus-nats-exporter:0.3.0 @@ -154,16 +156,23 @@ services: depends_on: - jaeger jaeger: - image: jaegertracing/all-in-one:1.42 + image: jaegertracing/all-in-one:1.49 container_name: jaeger hostname: jaeger restart: always networks: - main - ports: [ "6831:6831/udp", "6832:6832/udp", "5778:5778", "16686:16686", "14317:4317", "14318:4318", "14250:14250", "14268:14268", "14269:14269", "9411:9411" ] + ports: [ "6831:6831/udp", "6832:6832/udp", "5778:5778", "16686:16686", "14317:4317", "14318:4318", "14250:14250", "14268:14268", "14269:14269", "9411:9411", "16687:16687" ] + command: --query.ui-config /etc/jaeger/jaeger-ui.json environment: - - COLLECTOR_ZIPKIN_HOST_PORT=:9411 - COLLECTOR_OTLP_ENABLED=true + - METRICS_STORAGE_TYPE=prometheus + - PROMETHEUS_SERVER_URL=http://prometheus:9090 + - PROMETHEUS_QUERY_SUPPORT_SPANMETRICS_CONNECTOR=true + - PROMETHEUS_QUERY_NAMESPACE=span_metrics + - PROMETHEUS_QUERY_DURATION_UNIT=s + volumes: + - "./jaeger-ui.json:/etc/jaeger/jaeger-ui.json" nodeexporter: image: prom/node-exporter:v0.18.1 volumes: diff --git a/examples/tracing_middleware_compose/jaeger-ui.json b/examples/tracing_middleware_compose/jaeger-ui.json new file mode 100644 index 0000000..ce95752 --- /dev/null +++ b/examples/tracing_middleware_compose/jaeger-ui.json @@ -0,0 +1,8 @@ +{ + "monitor": { + "menuEnabled": true + }, + "dependencies": { + "menuEnabled": true + } +} \ No newline at end of file diff --git a/examples/tracing_middleware_compose/open_telemetry_config.yaml b/examples/tracing_middleware_compose/open_telemetry_config.yaml index 4f7077e..7c218d3 100644 --- a/examples/tracing_middleware_compose/open_telemetry_config.yaml +++ b/examples/tracing_middleware_compose/open_telemetry_config.yaml @@ -6,52 +6,31 @@ receivers: include: [ /var/log/*.log ] processors: -# attributes: -# actions: -# - action: insert -# key: log_file_name -# from_attribute: log.file.name -# - action: insert -# key: loki.attribute.labels -# value: log_file_name batch: + spanmetrics: + metrics_exporter: prometheus exporters: prometheus: namespace: test-space - const_labels: - label1: value1 - 'another label': spaced value send_timestamps: true metric_expiration: 180m resource_to_telemetry_conversion: enabled: true endpoint: "0.0.0.0:8889" -# loki: -# endpoint: "http://loki:3100/loki/api/v1/push" -# tenant_id: "example1" -# labels: -# attributes: -# log.file.name: "filename" -# container_name: "" -# container_id: "" jaeger: endpoint: jaeger:14250 tls: insecure: true extensions: health_check: - pprof: - endpoint: :1888 - zpages: - endpoint: :55679 service: - extensions: [pprof, zpages, health_check] + extensions: [health_check] pipelines: traces: receivers: [otlp] - processors: [batch] + processors: [batch, spanmetrics] exporters: [jaeger] metrics: receivers: [otlp] diff --git a/examples/tracing_middleware_compose/prometheus.yml b/examples/tracing_middleware_compose/prometheus.yml index 9646348..41dac18 100644 --- a/examples/tracing_middleware_compose/prometheus.yml +++ b/examples/tracing_middleware_compose/prometheus.yml @@ -13,7 +13,7 @@ scrape_configs: - job_name: 'otel-collector' scrape_interval: 10s static_configs: - - targets: [ 'otel-collector:8889' ] + - targets: [ 'open_telemetry:8889' ] - job_name: 'nodeexporter' static_configs: @@ -39,6 +39,11 @@ scrape_configs: static_configs: - targets: [ 'prometheus-nats-exporter:7777' ] + - job_name: 'jaeger' + metrics_path: /metrics + static_configs: + - targets: [ 'jaeger:14269' ] + alerting: alertmanagers: diff --git a/examples/tracing_middleware_compose/requirements.txt b/examples/tracing_middleware_compose/requirements.txt index 915c693..0f96236 100644 --- a/examples/tracing_middleware_compose/requirements.txt +++ b/examples/tracing_middleware_compose/requirements.txt @@ -1,4 +1,6 @@ -git+https://github.com/i2gor87/panini.git@feature/tracing-middleware +--index-url https://pypiserver.apps.colibriql-cluster.colibriql-infra.com +panini[tracing]==0.8.3b2 +panini==0.8.3b2 opentelemetry-api==1.19.0 opentelemetry-sdk==1.19.0 opentelemetry-exporter-otlp-proto-grpc==1.19.0 diff --git a/panini/middleware/tracing_middleware.py b/panini/middleware/tracing_middleware.py index 93d9299..0528372 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -1,6 +1,7 @@ try: from opentelemetry import trace - from opentelemetry.sdk.trace import TracerProvider, Tracer + from opentelemetry.trace import SpanKind + from opentelemetry.sdk.trace import TracerProvider, Tracer, Span from opentelemetry.sdk.trace.export import BatchSpanProcessor from opentelemetry.sdk.resources import SERVICE_NAME, Resource from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter @@ -189,7 +190,7 @@ async def trace_send_any( use_tracing = kwargs.pop("use_tracing", True) action = kwargs.pop("nats_action") existing_events: List[TracingEvent] = kwargs.pop("tracing_events", []) - if kwargs.get("use_current_span", False): + if kwargs.pop("use_current_span", False): ctx = trace.get_current_span().get_span_context() link = [trace.Link(ctx)] else: @@ -201,7 +202,7 @@ async def trace_send_any( if use_tracing is True: self.tracer.start_span(name=span_config.span_name) with self.tracer.start_as_current_span( - span_config.span_name, links=link, context=context + span_config.span_name, links=link, context=context, kind=SpanKind.SERVER ) as span: for attr_key, attr_value in span_config.span_attributes.items(): span.set_attribute(attr_key, attr_value) @@ -306,7 +307,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): span_attributes={}, ) with self.tracer.start_as_current_span( - span_config.span_name, context=context + span_config.span_name, context=context, kind=SpanKind.CLIENT ) as span: for ( attr_key, @@ -323,7 +324,7 @@ async def trace_listen_any(self, msg: Msg, callback, nats_action=None): ) try: response = await callback(msg) - if "tracing_events" in response.keys(): + if response and "tracing_events" in response.keys(): tracing_events: List[TracingEvent] = response.pop( "tracing_events", [] ) diff --git a/setup.py b/setup.py index d28cd96..e3adcb6 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ setup( name="panini", - version="0.8.3a11", + version="0.8.3b2", description="A python messaging framework for microservices based on NATS", long_description=long_description, long_description_content_type="text/x-rst",