Skip to content
This repository was archived by the owner on Feb 7, 2025. It is now read-only.
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion examples/tracing_middleware_compose/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
21 changes: 18 additions & 3 deletions examples/tracing_middleware_compose/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,11 @@ services:
- main
depends_on:
- natsserver
- jaeger
- grafana
- receiver
- prometheus
- open_telemetry
receiver:
build:
context: .
Expand All @@ -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"'
Expand Down Expand Up @@ -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:
Expand Down
8 changes: 8 additions & 0 deletions examples/tracing_middleware_compose/jaeger-ui.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
{
"monitor": {
"menuEnabled": true
},
"dependencies": {
"menuEnabled": true
}
}
29 changes: 4 additions & 25 deletions examples/tracing_middleware_compose/open_telemetry_config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
7 changes: 6 additions & 1 deletion examples/tracing_middleware_compose/prometheus.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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:
Expand Down
24 changes: 23 additions & 1 deletion examples/tracing_middleware_compose/receiver/app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"],
Expand All @@ -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}
Expand Down
12 changes: 7 additions & 5 deletions examples/tracing_middleware_compose/requirements.txt
Original file line number Diff line number Diff line change
@@ -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
PyYAML==6.0.1
20 changes: 19 additions & 1 deletion examples/tracing_middleware_compose/sender/app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"],
Expand Down Expand Up @@ -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('.')
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
service_name: abcdef
exporter_config:
endpoint: localhost:4317
endpoint: open_telemetry:4317
insecure: yes
headers: null
timeout: null
Expand Down
8 changes: 5 additions & 3 deletions panini/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
"""
Expand All @@ -148,15 +148,17 @@ 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(
routes=self.http,
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):
Expand Down
9 changes: 6 additions & 3 deletions panini/http_server/http_server_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {}
Expand All @@ -27,13 +28,15 @@ 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()

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
)
12 changes: 6 additions & 6 deletions panini/managers/middleware_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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)
Expand All @@ -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)
Expand All @@ -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

Expand Down
Loading