Skip to content
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
52 changes: 8 additions & 44 deletions gunicorn.conf.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
4 changes: 4 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
4 changes: 4 additions & 0 deletions requirements.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
135 changes: 99 additions & 36 deletions simone/dispatcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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)

Expand All @@ -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):
Expand Down Expand Up @@ -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
)
Expand All @@ -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:
Expand All @@ -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)
Expand Down
Loading
Loading