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/docker-compose.yml b/examples/tracing_middleware_compose/docker-compose.yml index 88f8042..268f95f 100644 --- a/examples/tracing_middleware_compose/docker-compose.yml +++ b/examples/tracing_middleware_compose/docker-compose.yml @@ -26,7 +26,11 @@ services: - main depends_on: - natsserver + - jaeger + - grafana - receiver + - prometheus + - open_telemetry receiver: build: context: . @@ -43,6 +47,10 @@ services: - main depends_on: - natsserver + - jaeger + - grafana + - prometheus + - open_telemetry nats-exporter: image: synadia/prometheus-nats-exporter:0.3.0 command: '-varz "http://nats-server:8222"' @@ -148,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/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/requirements.txt b/examples/tracing_middleware_compose/requirements.txt index c7d4d9b..0f96236 100644 --- a/examples/tracing_middleware_compose/requirements.txt +++ b/examples/tracing_middleware_compose/requirements.txt @@ -1,7 +1,9 @@ -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 +--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 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/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/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/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/managers/middleware_manager.py b/panini/managers/middleware_manager.py index 9677cd3..361d8dd 100644 --- a/panini/managers/middleware_manager.py +++ b/panini/managers/middleware_manager.py @@ -86,7 +86,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 +102,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 +116,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 +138,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..0528372 100644 --- a/panini/middleware/tracing_middleware.py +++ b/panini/middleware/tracing_middleware.py @@ -1,47 +1,95 @@ +try: + from opentelemetry import trace + 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 + from opentelemetry.trace.propagation.tracecontext import ( + TraceContextTextMapPropagator, + ) +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 typing import Optional, List from nats.aio.msg import Msg from panini.app import get_app -from opentelemetry import trace -from dataclasses import dataclass +from dataclasses import dataclass, field 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 +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] +@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): + @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 @@ -51,6 +99,28 @@ def decorated_function(*args, **kwargs): # the decorated function 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"] @@ -61,9 +131,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) @@ -72,68 +140,220 @@ def create_tracer(self): class TracingMiddleware(Middleware): - def __init__( - self, - tracing_config: dict, - **kwargs - ): - self._otel_tracer = OTELTracer(tracing_config=tracing_config, **kwargs) + """ + 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): + 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 - async def send_any(self, subject: str, message: Msg, send_func, *args, **kwargs): + def _get_span_name(self, action: str, subject: str) -> str: + return f"{action.upper()} {subject}" + + @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 + ): + verbose = kwargs.pop("verbose", False) carrier = {} - headers = {} + context = kwargs.pop("tracing_span_carrier", None) span_config = kwargs.get("span_config") - use_tracing = kwargs.get("use_tracing", True) - if kwargs.get("use_current_span", False): + use_tracing = kwargs.pop("use_tracing", True) + action = kwargs.pop("nats_action") + existing_events: List[TracingEvent] = kwargs.pop("tracing_events", []) + if kwargs.pop("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: + 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( + 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) + span.add_event( + action, + { + "nats.subject": subject, + "nats.message": json.dumps(message) if verbose else json.dumps(message)[:300], + }, + ) + for existing_event in existing_events: + span.add_event(existing_event.event_name, existing_event.event_data) + 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), } - if "use_tracing" in kwargs: - del kwargs['use_tracing'] - response = await send_func(subject, message, headers=headers) + kwargs.update({"headers": headers}) + try: + response = await send_func(subject, message, *args, **kwargs) + if response and verbose: + span.add_event("request_response", {"nats.message": response}) + return response + except Exception as exc: + span.record_exception(exc) + raise exc + response = await send_func(subject, message, *args, **kwargs) return response - async def listen_any(self, msg: Msg, callback): + 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 + ) + 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 + ) + 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 trace_listen_any(self, msg: Msg, callback, nats_action=None): 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) + 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: + verbose = listen_object._meta.get("verbose", False) + 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._get_span_name(nats_action, subject), + span_attributes={}, + ) + with self.tracer.start_as_current_span( + span_config.span_name, context=context, kind=SpanKind.CLIENT + ) 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( + nats_action, + { + "nats.subject": subject, + "nats.message": json.dumps(msg.data) if verbose else json.dumps(msg.data)[:300], + }, + ) + try: + response = await callback(msg) + if response and "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 + ) + if verbose: + 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!" + ) response = await callback(msg) - return response - return 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" + ) + return response + + async def listen_request(self, msg, callback): + response = await self.trace_listen_any( + msg, callback, nats_action="listen_request" + ) + return response 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..1b11281 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.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 cdcbf36..e3adcb6 100644 --- a/setup.py +++ b/setup.py @@ -17,9 +17,16 @@ """ +tracing_dependencies = [ + "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.1", + 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", @@ -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 + } )