diff --git a/gunicorn.conf.py b/gunicorn.conf.py index 09680b3..ac54ac3 100644 --- a/gunicorn.conf.py +++ b/gunicorn.conf.py @@ -50,50 +50,14 @@ def post_fork(server, worker): thread's state does. A TracerProvider built in the master would leave every worker holding a processor whose export thread only ever existed in the parent, silently exporting nothing. post_fork runs inside each freshly - forked worker, before it imports simone.wsgi, so DjangoInstrumentor is in - place before Django's own URL resolution and middleware load. + forked worker, before it imports simone.wsgi, so the instrumentation is in + place before Django's own machinery loads -- and, for the DB spans, before + Django opens its first connection. - Spans go straight to Tempo's OTLP/HTTP receiver (OTEL_EXPORTER_OTLP_ENDPOINT - in the compose environment), not through logit -- there is nothing logit - would add to spans an SDK already produced. Each one is a child of the span - logit lifts from nginx's access log line for the same request: nginx sets a - traceparent header and the default W3C propagator picks it up here with no - code of our own. - - A no-op when OTEL_EXPORTER_OTLP_ENDPOINT is unset, so running outside the - compose stack (script/run directly, the dev docker-compose.yml) doesn't - need a collector listening or spend every request's teardown waiting on a - connection refused. + The setup itself lives in simone/tracing.py so the management commands and + a dev server can opt in the same way. It's a no-op when + OTEL_EXPORTER_OTLP_ENDPOINT is unset. ''' - if not os.environ.get('OTEL_EXPORTER_OTLP_ENDPOINT'): - return - - from opentelemetry import trace - from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( - OTLPSpanExporter, - ) - from opentelemetry.instrumentation.django import DjangoInstrumentor - from opentelemetry.instrumentation.logging import LoggingInstrumentor - from opentelemetry.sdk.resources import Resource - from opentelemetry.sdk.trace import TracerProvider - from opentelemetry.sdk.trace.export import BatchSpanProcessor - - # service.name is what lets Tempo resolve a root service for these spans, - # matching the `set` component logit stamps onto nginx's own. - resource = Resource.create( - { - 'service.name': os.environ.get('OTEL_SERVICE_NAME', 'simone'), - 'service.namespace': 'xormedia', - } - ) - provider = TracerProvider(resource=resource) - # The exporter reads OTEL_EXPORTER_OTLP_ENDPOINT itself and appends - # /v1/traces per the OTLP spec. This package only speaks protobuf, which - # Tempo's receiver accepts natively. - provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter())) - trace.set_tracer_provider(provider) + from simone.tracing import configure_tracing - DjangoInstrumentor().instrument() - # Puts the active trace/span id into every log record, so a log line can be - # matched back to the trace it happened in. - LoggingInstrumentor().instrument(set_logging_format=True) + configure_tracing() diff --git a/pyproject.toml b/pyproject.toml index df16674..e464ba3 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -18,8 +18,12 @@ dependencies = [ "nltk>=3.6.6", "opentelemetry-api", "opentelemetry-exporter-otlp-proto-http", + "opentelemetry-instrumentation-dbapi", "opentelemetry-instrumentation-django", "opentelemetry-instrumentation-logging", + "opentelemetry-instrumentation-requests", + "opentelemetry-instrumentation-threading", + "opentelemetry-instrumentation-urllib", "opentelemetry-sdk", "pylev>=1.4.0", "requests>=2.26.0", diff --git a/requirements.txt b/requirements.txt index 2b08c99..4235202 100644 --- a/requirements.txt +++ b/requirements.txt @@ -33,8 +33,12 @@ opentelemetry-api==1.44.0 opentelemetry-exporter-otlp-proto-common==1.44.0 opentelemetry-exporter-otlp-proto-http==1.44.0 opentelemetry-instrumentation==0.65b0 +opentelemetry-instrumentation-dbapi==0.65b0 opentelemetry-instrumentation-django==0.65b0 opentelemetry-instrumentation-logging==0.65b0 +opentelemetry-instrumentation-requests==0.65b0 +opentelemetry-instrumentation-threading==0.65b0 +opentelemetry-instrumentation-urllib==0.65b0 opentelemetry-instrumentation-wsgi==0.65b0 opentelemetry-proto==1.44.0 opentelemetry-sdk==1.44.0 diff --git a/simone/dispatcher.py b/simone/dispatcher.py index c40ae78..fc0ce00 100644 --- a/simone/dispatcher.py +++ b/simone/dispatcher.py @@ -4,6 +4,9 @@ from django.conf import settings from django.db import close_old_connections, transaction from functools import wraps +from opentelemetry import context as otel_context, trace +from opentelemetry.context import Context +from opentelemetry.trace import SpanKind, Status, StatusCode from io import StringIO from logging import getLogger from os import environ, path @@ -25,25 +28,38 @@ max_workers=max_dispatchers, thread_name_prefix='simone-worker' ) +# No-op unless simone/tracing.py configured a real provider, so everything +# below is safe with tracing turned off -- the API's default provider hands +# back non-recording spans. +tracer = trace.get_tracer('simone.dispatcher') + def dispatch_with_error_reporting(func): @wraps(func) def wrap(self, context, *args, **kwargs): ret = None - with transaction.atomic(): - try: - ret = func(self, context, *args, **kwargs) - except Exception: - self.log.exception( - 'dispatch failed: context=%s, args=%s, kwargs=%s', - context, - args, - kwargs, - ) - context.say( - 'An error occured while responding to this message', - reply=True, - ) + # Span outside the transaction so it covers the commit too, which is + # where a slow write actually shows up. Recording the exception is the + # point rather than a bonus: this decorator deliberately swallows + # failures, so without it a broken handler is indistinguishable from a + # successful request. + with tracer.start_as_current_span(f'dispatch.{func.__name__}') as span: + with transaction.atomic(): + try: + ret = func(self, context, *args, **kwargs) + except Exception as e: + span.record_exception(e) + span.set_status(Status(StatusCode.ERROR)) + self.log.exception( + 'dispatch failed: context=%s, args=%s, kwargs=%s', + context, + args, + kwargs, + ) + context.say( + 'An error occured while responding to this message', + reply=True, + ) return ret return wrap @@ -55,16 +71,21 @@ def dispatch(func): @wraps(func) def wrap(self, context, *args, **kwargs): ret = None - with transaction.atomic(): - try: - ret = func(self, context, *args, **kwargs) - except Exception: - self.log.exception( - 'dispatch failed: context=%s, args=%s, kwargs=%s', - context, - args, - kwargs, - ) + # See dispatch_with_error_reporting above for why the span wraps the + # transaction and why the exception is recorded explicitly. + with tracer.start_as_current_span(f'dispatch.{func.__name__}') as span: + with transaction.atomic(): + try: + ret = func(self, context, *args, **kwargs) + except Exception as e: + span.record_exception(e) + span.set_status(Status(StatusCode.ERROR)) + self.log.exception( + 'dispatch failed: context=%s, args=%s, kwargs=%s', + context, + args, + kwargs, + ) return ret return wrap @@ -266,9 +287,19 @@ def command(self, context, text, **kwargs): self.log.debug('command: text=%s, kwargs=%s', text, kwargs) command_words, handler, command, text = self.find_command_handler(text) if handler: - handler.command( - context, command=command, text=text, dispatcher=self, **kwargs - ) + # Named for the command, so a trace says *which* command was slow + # rather than just "a command was". The handler class goes on as an + # attribute rather than into the name, to keep the name low + # cardinality. + with tracer.start_as_current_span(f'command.{command}') as span: + span.set_attribute('simone.handler', handler.__class__.__name__) + handler.command( + context, + command=command, + text=text, + dispatcher=self, + **kwargs, + ) else: self._did_you_mean(context, command_words) @@ -287,8 +318,14 @@ def left(self, *args, **kwargs): @dispatch def message(self, *args, **kwargs): + # Every registered message handler runs on every message, and several + # do real work per message -- handler_loud's count()+offset fetch, + # handler_responder's NLTK tokenize, handler_sparkles' read-modify- + # write. One span each is what makes it obvious which one costs. for handler in self.messages: - handler.message(*args, dispatcher=self, **kwargs) + name = handler.__class__.__name__ + with tracer.start_as_current_span(f'message.{name}'): + handler.message(*args, dispatcher=self, **kwargs) @dispatch def removed(self, *args, **kwargs): @@ -356,8 +393,15 @@ def tick(self, now): bot_user_id=workspace.bot_user_id, channel=channel, ) - handler.cron(context, cron=cron, dispatcher=self) - except Exception: + name = handler.__class__.__name__ + with tracer.start_as_current_span(f'cron.{name}') as span: + span.set_attribute('simone.cron.when', cron['when']) + span.set_attribute( + 'simone.cron.channel', cron['channel'] + ) + handler.cron(context, cron=cron, dispatcher=self) + except Exception as e: + trace.get_current_span().record_exception(e) self.log.exception( 'tick: cron=%s failed for workspace=%s', cron, workspace ) @@ -372,6 +416,15 @@ def __init__(self, dispatcher): def run(self): self.log.info('run: starting') + # Detach from whatever context started this thread. ThreadingInstrumentor + # (simone/tracing.py) propagates the spawning context into new threads, + # which is exactly what we want for slack_bolt's listener executor and + # wrong here: this thread lives for the life of the process, so anything + # it inherited would parent every tick, forever, to one long-finished + # span. Today Cron is started from simone/wsgi.py at import, when no + # span is active, so this changes nothing -- it's here so that stays + # true if it's ever started from somewhere else. + otel_context.attach(Context()) self.stopper = Event() running = True while running: @@ -382,13 +435,23 @@ def run(self): # the thread as well as wrap each time around in calls to check # our database connections health (name doesn't match # functionality) - close_old_connections() - try: - self.dispatcher.tick(datetime.utcnow()) - except Exception: - self.log.exception('run: tick failed') - finally: + # + # A root span per tick -- its own trace, since nothing requested + # it. SERVER-ish work with no caller, so INTERNAL is the honest + # kind. close_old_connections is inside it because reconnect cost + # is part of what a tick spends. + with tracer.start_as_current_span( + 'cron.tick', kind=SpanKind.INTERNAL + ) as span: close_old_connections() + try: + self.dispatcher.tick(datetime.utcnow()) + except Exception as e: + span.record_exception(e) + span.set_status(Status(StatusCode.ERROR)) + self.log.exception('run: tick failed') + finally: + close_old_connections() elapsed = time() - start pause = 60 - elapsed self.log.debug('run: elapsed=%f, pause=%f', elapsed, pause) diff --git a/simone/tracing.py b/simone/tracing.py new file mode 100644 index 0000000..c58a64b --- /dev/null +++ b/simone/tracing.py @@ -0,0 +1,126 @@ +''' +OpenTelemetry setup, in one place so anything that runs the app can opt in -- +gunicorn (see gunicorn.conf.py's post_fork), the management commands, or a +dev server. Spans go to whatever OTEL_EXPORTER_OTLP_ENDPOINT names; in the +compose stack that's Grafana Tempo's OTLP/HTTP receiver. +''' + +from logging import getLogger +from os import environ + +log = getLogger('tracing') + + +def configure_tracing(): + ''' + Build a TracerProvider and install every instrumentor we use. + + A no-op when OTEL_EXPORTER_OTLP_ENDPOINT is unset, so running outside the + compose stack doesn't need a collector listening or spend every request's + teardown retrying a connection refused. + + WHERE this is called from matters, and isn't arbitrary: + + * The BatchSpanProcessor below owns a background export thread, and a + thread does not survive fork(). Under gunicorn this has to run in the + worker, post-fork, or the provider's exporter only ever exists in the + master and nothing is ever sent. + + * trace_integration() patches mysql.connector.connect, so it only + affects connections opened afterwards. Django opens its first + connection when simone.wsgi is imported, which is after post_fork -- + so this is correct today, but it's the reason this can't drift into + an AppConfig.ready() or a middleware later. + + Returns True if tracing was configured, False if it was skipped. + ''' + if not environ.get('OTEL_EXPORTER_OTLP_ENDPOINT'): + return False + + from opentelemetry import trace + from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( + OTLPSpanExporter, + ) + from opentelemetry.instrumentation.dbapi import trace_integration + from opentelemetry.instrumentation.django import DjangoInstrumentor + from opentelemetry.instrumentation.logging import LoggingInstrumentor + from opentelemetry.instrumentation.requests import RequestsInstrumentor + from opentelemetry.instrumentation.threading import ThreadingInstrumentor + from opentelemetry.instrumentation.urllib import URLLibInstrumentor + from opentelemetry.sdk.resources import Resource + from opentelemetry.sdk.trace import TracerProvider + from opentelemetry.sdk.trace.export import BatchSpanProcessor + + # service.name is what lets Tempo resolve a root service for these spans, + # matching the `set` component logit stamps onto nginx's own. + resource = Resource.create( + { + 'service.name': environ.get('OTEL_SERVICE_NAME', 'simone'), + 'service.namespace': 'xormedia', + } + ) + provider = TracerProvider(resource=resource) + # The exporter reads OTEL_EXPORTER_OTLP_ENDPOINT itself and appends + # /v1/traces per the OTLP spec. This package only speaks protobuf, which + # Tempo's receiver accepts natively. + provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter())) + trace.set_tracer_provider(provider) + + # One SERVER span per request. Note it produces no DB spans of its own -- + # that's what trace_integration below is for. + DjangoInstrumentor().instrument() + # Puts the active trace/span id into every log record, so a log line can + # be matched back to the trace it happened in. + LoggingInstrumentor().instrument(set_logging_format=True) + + # Outbound HTTP, which takes two instrumentors because the app makes its + # calls two different ways: + # requests -- the shared Session in simone/handlers.py, used by every + # handler/chat API call (weather, jokes, quotes, images, stonks). + # urllib -- slack_sdk's sync WebClient, which builds an OpenerDirector + # over urllib.request rather than using requests or urllib3. Every + # chat_postMessage/reactions_add/conversations_info goes through here, + # so without this the Slack egress on nearly every request is + # invisible. + RequestsInstrumentor().instrument() + URLLibInstrumentor().instrument() + + # Emits no telemetry itself; it propagates the active context into threads + # and ThreadPoolExecutor workers. Load-bearing here: slack_bolt runs its + # listeners on the executor in simone/dispatcher.py, so handler work + # happens *after* the HTTP response is sent (Slack wants an ack in ~3s). + # Without this every handler span would be an orphan root instead of a + # child of the request that caused it. A child outliving its parent is + # normal for async work and renders correctly in Tempo. + # + # It also covers the private ThreadPoolExecutor in handler/chat/stonks.py. + # Cron opts back out explicitly -- see Cron.run in simone/dispatcher.py. + ThreadingInstrumentor().instrument() + + # DB spans, via the dbapi instrumentation directly rather than + # opentelemetry-instrumentation-mysql. That package is a single + # wrap_connect call wearing a version bound of + # "mysql-connector-python >= 8.0, < 10.0" -- we're on 26.x, so its + # instrument() would raise DependencyConflict, log an error and return + # WITHOUT instrumenting. It would look installed and quietly do nothing. + # Calling the underlying function skips the bogus gate and a dependency. + # + # This reaches the ORM because mysql/connector/django/base.py does + # `cnx = mysql.connector.connect(...)` -- a live module attribute lookup, + # so the patch intercepts it and Django's connection is the traced proxy. + # + # Imported here rather than at module scope: dev settings fall back to + # sqlite3 when SIMONE_DB_NAME is unset, and this shouldn't be a hard + # requirement there. + # + # No enable_commenter: sqlcommenter rewrites each query with a unique + # traceparent comment, which defeats statement caching, and upstream warns + # it's pathological with prepared cursors. + try: + import mysql.connector + + trace_integration(mysql.connector, 'connect', 'mysql') + except ImportError: + log.info('configure_tracing: mysql.connector absent, no DB spans') + + return True