mirror of
https://github.com/simonw/datasette.git
synced 2026-09-08 09:34:05 +02:00
Compare commits
No commits in common. "asg017/otel-phase1-6-plugin-kit" and "main" have entirely different histories.
asg017/ote
...
main
21 changed files with 156 additions and 6722 deletions
159
datasette/app.py
159
datasette/app.py
|
|
@ -49,16 +49,6 @@ from .events import Event
|
|||
from .plugins import DEFAULT_PLUGINS, get_plugins, pm
|
||||
from .renderer import json_renderer
|
||||
from .resources import DatabaseResource, TableResource
|
||||
from .telemetry import (
|
||||
TelemetryMiddleware,
|
||||
_in_datasette_client,
|
||||
clamp_http_method,
|
||||
register_datasette,
|
||||
request_span,
|
||||
tracer,
|
||||
unregister_datasette,
|
||||
)
|
||||
from .telemetry_registry import HTTP_ROUTE, STARTUP
|
||||
from .tokens import TokenInvalid
|
||||
from .tracer import AsgiTracer
|
||||
from .url_builder import Urls
|
||||
|
|
@ -174,9 +164,8 @@ app_root = Path(__file__).parent.parent
|
|||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
# _in_datasette_client itself lives in telemetry.py so the request span
|
||||
# middleware can read it without a circular import; its writers
|
||||
# (_DatasetteClientContext) and reader (in_client()) both live here.
|
||||
# Context variable to track when code is executing within a datasette.client request
|
||||
_in_datasette_client = contextvars.ContextVar("in_datasette_client", default=False)
|
||||
|
||||
|
||||
class _DatasetteClientContext:
|
||||
|
|
@ -649,10 +638,6 @@ class Datasette:
|
|||
self.root_enabled = False
|
||||
self.default_deny = default_deny
|
||||
self.client = DatasetteClient(self)
|
||||
# Last, so that the observable-gauge callbacks - which may fire on the
|
||||
# SDK's collection thread the instant this returns - never see a
|
||||
# half-built instance.
|
||||
register_datasette(self)
|
||||
|
||||
async def apply_metadata_json(self):
|
||||
# Apply any metadata entries from metadata.json to the internal tables
|
||||
|
|
@ -793,73 +778,57 @@ class Datasette:
|
|||
# This must be called for Datasette to be in a usable state
|
||||
if self._startup_invoked:
|
||||
return
|
||||
# `datasette serve` calls invoke_startup() before uvicorn starts, so
|
||||
# on the CLI path every span its children create - the register_*
|
||||
# hook dispatches, the internal catalog's db.query/db.write spans,
|
||||
# and the prepare_connection warm-up of the read connections those
|
||||
# touch - would otherwise be its own orphan root trace: around twenty
|
||||
# of them on a fresh instance. Bracketing the whole thing gives them
|
||||
# somewhere to belong. An ASGI-hosted or programmatic deployment
|
||||
# reaches here instead through AsgiRunOnFirstRequest, in which case
|
||||
# this span nests under the first request's own span - honest enough,
|
||||
# since it genuinely is that request's latency.
|
||||
# A connection warmed lazily later, by a request touching a new
|
||||
# database for the first time, nests under that request instead:
|
||||
# this span has already ended by then.
|
||||
with tracer.start_as_current_span(STARTUP):
|
||||
# Register event classes
|
||||
event_classes = []
|
||||
for hook in pm.hook.register_events(datasette=self):
|
||||
extra_classes = await await_me_maybe(hook)
|
||||
if extra_classes:
|
||||
event_classes.extend(extra_classes)
|
||||
self.event_classes = tuple(event_classes)
|
||||
# Register event classes
|
||||
event_classes = []
|
||||
for hook in pm.hook.register_events(datasette=self):
|
||||
extra_classes = await await_me_maybe(hook)
|
||||
if extra_classes:
|
||||
event_classes.extend(extra_classes)
|
||||
self.event_classes = tuple(event_classes)
|
||||
|
||||
# Register actions, but watch out for duplicate name/abbr
|
||||
action_names = {}
|
||||
action_abbrs = {}
|
||||
for hook in pm.hook.register_actions(datasette=self):
|
||||
if hook:
|
||||
for action in hook:
|
||||
if (
|
||||
action.name in action_names
|
||||
and action != action_names[action.name]
|
||||
):
|
||||
raise StartupError(f"Duplicate action name: {action.name}")
|
||||
if (
|
||||
action.abbr
|
||||
and action.abbr in action_abbrs
|
||||
and action != action_abbrs[action.abbr]
|
||||
):
|
||||
raise StartupError(f"Duplicate action abbr: {action.abbr}")
|
||||
action_names[action.name] = action
|
||||
if action.abbr:
|
||||
action_abbrs[action.abbr] = action
|
||||
self.actions[action.name] = action
|
||||
# Register actions, but watch out for duplicate name/abbr
|
||||
action_names = {}
|
||||
action_abbrs = {}
|
||||
for hook in pm.hook.register_actions(datasette=self):
|
||||
if hook:
|
||||
for action in hook:
|
||||
if (
|
||||
action.name in action_names
|
||||
and action != action_names[action.name]
|
||||
):
|
||||
raise StartupError(f"Duplicate action name: {action.name}")
|
||||
if (
|
||||
action.abbr
|
||||
and action.abbr in action_abbrs
|
||||
and action != action_abbrs[action.abbr]
|
||||
):
|
||||
raise StartupError(f"Duplicate action abbr: {action.abbr}")
|
||||
action_names[action.name] = action
|
||||
if action.abbr:
|
||||
action_abbrs[action.abbr] = action
|
||||
self.actions[action.name] = action
|
||||
|
||||
# Register column types (classes, not instances)
|
||||
self._column_types = {}
|
||||
for hook in pm.hook.register_column_types(datasette=self):
|
||||
if hook:
|
||||
for ct_cls in hook:
|
||||
if ct_cls.name in self._column_types:
|
||||
raise StartupError(
|
||||
f"Duplicate column type name: {ct_cls.name}"
|
||||
)
|
||||
self._column_types[ct_cls.name] = ct_cls
|
||||
# Register column types (classes, not instances)
|
||||
self._column_types = {}
|
||||
for hook in pm.hook.register_column_types(datasette=self):
|
||||
if hook:
|
||||
for ct_cls in hook:
|
||||
if ct_cls.name in self._column_types:
|
||||
raise StartupError(f"Duplicate column type name: {ct_cls.name}")
|
||||
self._column_types[ct_cls.name] = ct_cls
|
||||
|
||||
for hook in pm.hook.prepare_jinja2_environment(
|
||||
env=self._jinja_env, datasette=self
|
||||
):
|
||||
await await_me_maybe(hook)
|
||||
# Ensure internal tables and metadata are populated before startup hooks
|
||||
await self._refresh_schemas()
|
||||
await self._save_queries_from_config()
|
||||
# Load column_types from config into internal DB
|
||||
await self._apply_column_types_config()
|
||||
for hook in pm.hook.startup(datasette=self):
|
||||
await await_me_maybe(hook)
|
||||
self._startup_invoked = True
|
||||
for hook in pm.hook.prepare_jinja2_environment(
|
||||
env=self._jinja_env, datasette=self
|
||||
):
|
||||
await await_me_maybe(hook)
|
||||
# Ensure internal tables and metadata are populated before startup hooks
|
||||
await self._refresh_schemas()
|
||||
await self._save_queries_from_config()
|
||||
# Load column_types from config into internal DB
|
||||
await self._apply_column_types_config()
|
||||
for hook in pm.hook.startup(datasette=self):
|
||||
await await_me_maybe(hook)
|
||||
self._startup_invoked = True
|
||||
|
||||
def sign(self, value, namespace="default"):
|
||||
return URLSafeSerializer(self._secret, namespace).dumps(value)
|
||||
|
|
@ -992,10 +961,6 @@ class Datasette:
|
|||
if self._closed:
|
||||
return
|
||||
self._closed = True
|
||||
# Stop reporting gauges before tearing anything down, so a collection
|
||||
# cycle landing mid-close cannot observe a half-closed instance. The
|
||||
# WeakSet would drop it eventually anyway; this makes it immediate.
|
||||
unregister_datasette(self)
|
||||
first_exception = None
|
||||
dbs = list(self.databases.values()) + [self._internal_database]
|
||||
for db in dbs:
|
||||
|
|
@ -2889,12 +2854,6 @@ class Datasette:
|
|||
asgi = AsgiRunOnFirstRequest(asgi, on_startup=[self._startup_sequence])
|
||||
for wrapper in pm.hook.asgi_wrapper(datasette=self):
|
||||
asgi = wrapper(asgi)
|
||||
# Outermost, deliberately: plugin asgi_wrapper() middleware, the
|
||||
# CSRF layer and the first-request startup fallback all run *inside*
|
||||
# this span, so a span created by an instrumented plugin - or by
|
||||
# startup work triggered by the first request - parents to the
|
||||
# request instead of becoming its own orphan root trace.
|
||||
asgi = TelemetryMiddleware(asgi)
|
||||
return asgi
|
||||
|
||||
|
||||
|
|
@ -2975,26 +2934,8 @@ class DatasetteRouter:
|
|||
match, view = resolve_routes(self.routes, path)
|
||||
|
||||
if match is None:
|
||||
# No route matched, so the span keeps the bare method name it was
|
||||
# given at the edge and gets no http.route. That is what semantic
|
||||
# conventions ask for when the route is unknown.
|
||||
return await self.handle_404(request, send)
|
||||
|
||||
# The request span was started at the ASGI edge, before routing, so it
|
||||
# carries only the method as a name. Now that the route is known, give
|
||||
# it the `{method} {route}` shape semantic conventions want, and the
|
||||
# http.route attribute - the low-cardinality counterpart to url.path,
|
||||
# and so the one to group by.
|
||||
span = request_span(scope)
|
||||
if span is not None:
|
||||
route = match.re.pattern
|
||||
span.set_attribute(HTTP_ROUTE, route)
|
||||
# Clamped, for the same reason the middleware clamps it: the method
|
||||
# is a client-controlled string, and an unclamped one here would
|
||||
# put attacker-supplied text back into the span name that the
|
||||
# middleware just kept out of it.
|
||||
span.update_name(f"{clamp_http_method(request.method)} {route}")
|
||||
|
||||
new_scope = dict(scope, url_route={"kwargs": match.groupdict()})
|
||||
request.scope = new_scope
|
||||
try:
|
||||
|
|
|
|||
|
|
@ -157,11 +157,7 @@ async def inspect_(files, sqlite_extensions):
|
|||
app = Datasette([], immutables=files, sqlite_extensions=sqlite_extensions)
|
||||
data = {}
|
||||
for name, database in app.databases.items():
|
||||
|
||||
def _inspect_tables(conn):
|
||||
return inspect_tables(conn, {})
|
||||
|
||||
tables = await database.execute_fn(_inspect_tables)
|
||||
tables = await database.execute_fn(lambda conn: inspect_tables(conn, {}))
|
||||
data[name] = {
|
||||
"hash": database.hash,
|
||||
"size": database.size,
|
||||
|
|
|
|||
|
|
@ -1,55 +1,18 @@
|
|||
import asyncio
|
||||
import atexit
|
||||
import contextvars
|
||||
import inspect
|
||||
import os
|
||||
import queue
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
from collections import namedtuple
|
||||
from pathlib import Path
|
||||
|
||||
import sqlite_utils
|
||||
from opentelemetry import context as otel_context_api
|
||||
from opentelemetry.trace import Status, StatusCode
|
||||
|
||||
from .inspect import inspect_hash
|
||||
from .telemetry import (
|
||||
callback_name,
|
||||
linked_root_span_kwargs,
|
||||
record_operation_duration,
|
||||
record_query_interrupted,
|
||||
record_write_queue_wait,
|
||||
sql_attribute,
|
||||
sql_operation_name,
|
||||
tracer,
|
||||
)
|
||||
from .telemetry_registry import (
|
||||
CALLBACK,
|
||||
DB_COLLECTION_NAME,
|
||||
DB_NAMESPACE,
|
||||
DB_OPERATION_NAME,
|
||||
DB_QUERY,
|
||||
DB_QUERY_EXECUTE,
|
||||
DB_QUERY_TEXT,
|
||||
DB_SYSTEM,
|
||||
DB_WRITE_EXECUTE,
|
||||
DB_WRITE_QUEUE_WAIT,
|
||||
EXECUTEMANY,
|
||||
EXECUTESCRIPT,
|
||||
INTERRUPTED,
|
||||
ISOLATED_CONNECTION,
|
||||
PARAM_COUNT,
|
||||
PARAM_SETS,
|
||||
ROWS_RETURNED,
|
||||
SQL_ERROR_SUPPRESSED,
|
||||
TIME_LIMIT_MS,
|
||||
TRANSACTION,
|
||||
TRUNCATED,
|
||||
)
|
||||
from .tracer import trace
|
||||
from .utils import (
|
||||
call_with_supported_arguments,
|
||||
|
|
@ -294,25 +257,10 @@ class Database:
|
|||
cursor, return_all=return_all, returning_limit=returning_limit
|
||||
)
|
||||
|
||||
# SIM117 wants these two context managers merged. They are kept nested
|
||||
# deliberately: the hand-rolled tracer's wrapper is on its way out, and
|
||||
# nesting makes removing it a single-line deletion.
|
||||
with trace( # noqa: SIM117
|
||||
"sql", database=self.name, sql=sql.strip(), params=params
|
||||
):
|
||||
with tracer.start_as_current_span(DB_QUERY, kind=DB_QUERY.kind) as span:
|
||||
span.set_attribute(DB_SYSTEM, "sqlite")
|
||||
span.set_attribute(DB_NAMESPACE, self.name)
|
||||
span.set_attribute(DB_QUERY_TEXT, sql_attribute(sql))
|
||||
operation_name = sql_operation_name(sql)
|
||||
if operation_name:
|
||||
span.set_attribute(DB_OPERATION_NAME, operation_name)
|
||||
if params:
|
||||
span.set_attribute(PARAM_COUNT, len(params))
|
||||
with record_operation_duration(self.name, "write"):
|
||||
results = await self._execute_write_fn(
|
||||
_inner, block=block, request=request, transaction=transaction
|
||||
)
|
||||
with trace("sql", database=self.name, sql=sql.strip(), params=params):
|
||||
results = await self.execute_write_fn(
|
||||
_inner, block=block, request=request, transaction=transaction
|
||||
)
|
||||
return results
|
||||
|
||||
async def execute_write_script(self, sql, block=True, request=None):
|
||||
|
|
@ -321,23 +269,10 @@ class Database:
|
|||
def _inner(conn):
|
||||
return conn.executescript(sql)
|
||||
|
||||
# Nested on purpose - see the note in execute_write().
|
||||
with trace( # noqa: SIM117
|
||||
"sql", database=self.name, sql=sql.strip(), executescript=True
|
||||
):
|
||||
# No db.operation.name here, deliberately: executescript() runs
|
||||
# several semicolon-separated statements, and semantic conventions
|
||||
# say the attribute should not be extracted from query text that
|
||||
# can hold more than one operation - see sql_operation_name().
|
||||
with tracer.start_as_current_span(DB_QUERY, kind=DB_QUERY.kind) as span:
|
||||
span.set_attribute(DB_SYSTEM, "sqlite")
|
||||
span.set_attribute(DB_NAMESPACE, self.name)
|
||||
span.set_attribute(DB_QUERY_TEXT, sql_attribute(sql))
|
||||
span.set_attribute(EXECUTESCRIPT, True)
|
||||
with record_operation_duration(self.name, "write"):
|
||||
results = await self._execute_write_fn(
|
||||
_inner, block=block, transaction=False, request=request
|
||||
)
|
||||
with trace("sql", database=self.name, sql=sql.strip(), executescript=True):
|
||||
results = await self.execute_write_fn(
|
||||
_inner, block=block, transaction=False, request=request
|
||||
)
|
||||
return results
|
||||
|
||||
async def execute_write_many(self, sql, params_seq, block=True, request=None):
|
||||
|
|
@ -354,27 +289,12 @@ class Database:
|
|||
|
||||
return conn.executemany(sql, count_params(params_seq)), count
|
||||
|
||||
# Nested on purpose - see the note in execute_write().
|
||||
with trace(
|
||||
"sql", database=self.name, sql=sql.strip(), executemany=True
|
||||
) as kwargs:
|
||||
with tracer.start_as_current_span(DB_QUERY, kind=DB_QUERY.kind) as span:
|
||||
span.set_attribute(DB_SYSTEM, "sqlite")
|
||||
span.set_attribute(DB_NAMESPACE, self.name)
|
||||
span.set_attribute(DB_QUERY_TEXT, sql_attribute(sql))
|
||||
span.set_attribute(EXECUTEMANY, True)
|
||||
# A single statement run with many parameter sets, so unlike
|
||||
# execute_write_script() there is exactly one operation to name.
|
||||
operation_name = sql_operation_name(sql)
|
||||
if operation_name:
|
||||
span.set_attribute(DB_OPERATION_NAME, operation_name)
|
||||
with record_operation_duration(self.name, "write"):
|
||||
results, count = await self._execute_write_fn(
|
||||
_inner, block=block, request=request
|
||||
)
|
||||
# count is the number of parameter *sets* consumed by
|
||||
# executemany(), not a row count - executemany returns no rows.
|
||||
span.set_attribute(PARAM_SETS, count)
|
||||
results, count = await self.execute_write_fn(
|
||||
_inner, block=block, request=request
|
||||
)
|
||||
kwargs["count"] = count
|
||||
return results
|
||||
|
||||
|
|
@ -396,69 +316,26 @@ class Database:
|
|||
# Was probably a memory connection
|
||||
pass
|
||||
|
||||
# One db.query span here, like execute_fn() / execute_write_fn().
|
||||
# The wrap must NOT move into _send_to_write_thread(): that is the
|
||||
# shared tail for every write, and for block=False it is where the
|
||||
# link back to this span is captured - a span opened there would be
|
||||
# the link target for its own children.
|
||||
with tracer.start_as_current_span(DB_QUERY, kind=DB_QUERY.kind) as span:
|
||||
span.set_attribute(DB_SYSTEM, "sqlite")
|
||||
span.set_attribute(DB_NAMESPACE, self.name)
|
||||
span.set_attribute(CALLBACK, callback_name(fn))
|
||||
# "write" when mutable because the call blocks the write queue;
|
||||
# "read" when immutable, where it runs on the read pool.
|
||||
with record_operation_duration(self.name, "write" if write else "read"):
|
||||
if self.ds.executor is None:
|
||||
# non-threaded mode
|
||||
return _run()
|
||||
if not write:
|
||||
# Immutable database - no writes can ever occur, so there
|
||||
# is no write queue to block; run against a fresh
|
||||
# read-only connection. copy_context() carries the
|
||||
# caller's otel context onto the worker thread - see the
|
||||
# notes in _execute_fn() for why it must be a fresh copy
|
||||
# per submit and why carrying every ContextVar is safe.
|
||||
ctx = contextvars.copy_context()
|
||||
return await asyncio.get_running_loop().run_in_executor(
|
||||
self.ds.executor, ctx.run, _run
|
||||
)
|
||||
# Threaded mode - send to write thread
|
||||
return await self._send_to_write_thread(fn, isolated_connection=True)
|
||||
if self.ds.executor is None:
|
||||
# non-threaded mode
|
||||
return _run()
|
||||
if not write:
|
||||
# Immutable database - no writes can ever occur, so there is no
|
||||
# write queue to block; run against a fresh read-only connection
|
||||
return await asyncio.get_running_loop().run_in_executor(
|
||||
self.ds.executor, _run
|
||||
)
|
||||
# Threaded mode - send to write thread
|
||||
return await self._send_to_write_thread(fn, isolated_connection=True)
|
||||
|
||||
async def analyze_sql(self, sql, params=None) -> SQLAnalysis:
|
||||
self._check_not_closed()
|
||||
|
||||
def _analyze_sql(conn):
|
||||
return analyze_sql_tables(conn, sql, params, database_name=self.name)
|
||||
|
||||
return await self.execute_isolated_fn(_analyze_sql)
|
||||
return await self.execute_isolated_fn(
|
||||
lambda conn: analyze_sql_tables(conn, sql, params, database_name=self.name)
|
||||
)
|
||||
|
||||
async def execute_write_fn(self, fn, block=True, transaction=True, request=None):
|
||||
"""Run `fn(conn)` on the write connection, traced as one database call.
|
||||
|
||||
The public entry point for callback-style writes. Instrumented like
|
||||
`execute_write()`: one `db.query` span (with `datasette.callback` in
|
||||
place of `db.query.text`) above the `db.write.queue_wait` and
|
||||
`db.write.execute` spans the write thread emits. The SQL-string write
|
||||
methods call `_execute_write_fn()` directly, so they never get a
|
||||
second span. For `block=False` this span ends at enqueue and the
|
||||
write-thread spans become roots carrying a link back to it, exactly
|
||||
as for `execute_write(block=False)`.
|
||||
"""
|
||||
self._check_not_closed()
|
||||
# The raw fn's name, before _wrap_fn_with_hooks() replaces it with a
|
||||
# wrapper - otherwise every write would report the wrapper's name.
|
||||
name = callback_name(fn)
|
||||
with tracer.start_as_current_span(DB_QUERY, kind=DB_QUERY.kind) as span:
|
||||
span.set_attribute(DB_SYSTEM, "sqlite")
|
||||
span.set_attribute(DB_NAMESPACE, self.name)
|
||||
span.set_attribute(CALLBACK, name)
|
||||
with record_operation_duration(self.name, "write"):
|
||||
return await self._execute_write_fn(
|
||||
fn, block=block, transaction=transaction, request=request
|
||||
)
|
||||
|
||||
async def _execute_write_fn(self, fn, block=True, transaction=True, request=None):
|
||||
self._check_not_closed()
|
||||
pending_events = []
|
||||
|
||||
|
|
@ -551,21 +428,8 @@ class Database:
|
|||
task_id = uuid.uuid5(uuid.NAMESPACE_DNS, "datasette.io")
|
||||
loop = asyncio.get_running_loop()
|
||||
reply_future = loop.create_future()
|
||||
# The otel Context and enqueue timestamp are captured here, on the
|
||||
# event loop, for the db.write.queue_wait span built at dequeue time.
|
||||
# `block` travels too - it decides parent vs. link; see `_execute_writes`.
|
||||
self._write_queue.put(
|
||||
WriteTask(
|
||||
fn,
|
||||
task_id,
|
||||
loop,
|
||||
reply_future,
|
||||
isolated_connection,
|
||||
transaction,
|
||||
otel_context_api.get_current(),
|
||||
time.time_ns(),
|
||||
block,
|
||||
)
|
||||
WriteTask(fn, task_id, loop, reply_future, isolated_connection, transaction)
|
||||
)
|
||||
if block:
|
||||
return await reply_future
|
||||
|
|
@ -579,16 +443,6 @@ class Database:
|
|||
conn = None
|
||||
try:
|
||||
conn = self.connect(write=True)
|
||||
# This warm-up runs before any write has ever been queued, so
|
||||
# there is no captured caller context to attach - and a raw
|
||||
# threading.Thread does not inherit the context of whoever started
|
||||
# it. Spans created by plugin hooks here are therefore roots even
|
||||
# when the write thread is started from inside invoke_startup():
|
||||
# its datasette.startup span is current on the event loop but does
|
||||
# not cross this thread boundary. Read connections differ - they
|
||||
# warm up inside executor tasks submitted with copy_context(), so
|
||||
# their prepare_connection spans do nest under whoever triggered
|
||||
# them.
|
||||
self.ds._prepare_connection(conn, self.name)
|
||||
except Exception as e: # noqa: BLE001
|
||||
# Stored and re-raised to whoever queues the next write
|
||||
|
|
@ -603,133 +457,42 @@ class Database:
|
|||
# Best-effort close as the write thread exits
|
||||
pass
|
||||
return
|
||||
# `task.block` decides how this task's spans relate to the
|
||||
# context captured at enqueue time:
|
||||
#
|
||||
# - block=True: the caller genuinely awaits the reply, so
|
||||
# containment is accurate. Restore that context as current
|
||||
# (attach below) so db.write.queue_wait/db.write.execute parent
|
||||
# normally to the request that queued them. The token must be
|
||||
# detached below in `finally` - a leaked token silently
|
||||
# poisons this thread's ambient context for every write
|
||||
# processed after it, and a *wrong*-token detach only logs a
|
||||
# warning rather than raising, so this pairing is load-bearing
|
||||
# and easy to get wrong silently.
|
||||
# - block=False: the caller returned already without awaiting,
|
||||
# so the enqueueing span may have closed before this task's
|
||||
# spans even start. The enqueueing request *caused* this write
|
||||
# without *containing* it, so nothing is attached here -
|
||||
# linked_root_span_kwargs() makes each write span a root with
|
||||
# a Link back to the enqueueing span (see its docstring for
|
||||
# the full rationale), built once into `write_span_kwargs`
|
||||
# and spread into every start_span call below.
|
||||
token = None
|
||||
write_span_kwargs = {}
|
||||
if task.block:
|
||||
token = otel_context_api.attach(task.otel_context)
|
||||
exception = None
|
||||
result = None
|
||||
if conn_exception is not None:
|
||||
exception = conn_exception
|
||||
elif task.isolated_connection:
|
||||
try:
|
||||
isolated_connection = self.connect(write=True)
|
||||
try:
|
||||
result = task.fn(isolated_connection)
|
||||
finally:
|
||||
isolated_connection.close()
|
||||
try:
|
||||
self._all_file_connections.remove(isolated_connection)
|
||||
except ValueError:
|
||||
# Was probably a memory connection
|
||||
pass
|
||||
except Exception as e: # noqa: BLE001
|
||||
# Write thread must survive any task failure or the database wedges
|
||||
sys.stderr.write(f"{e}\n")
|
||||
sys.stderr.flush()
|
||||
exception = e
|
||||
else:
|
||||
write_span_kwargs = linked_root_span_kwargs(task.otel_context)
|
||||
try:
|
||||
exception = None
|
||||
result = None
|
||||
# Explicit start_time/end_time rather than a `with` block:
|
||||
# this span's duration is the time the task actually spent
|
||||
# waiting in the queue (enqueue -> dequeue), not the near-
|
||||
# zero time spent constructing/ending the span object here.
|
||||
dequeued_at_ns = time.time_ns()
|
||||
tracer.start_span(
|
||||
DB_WRITE_QUEUE_WAIT,
|
||||
start_time=task.enqueued_at_ns,
|
||||
**write_span_kwargs,
|
||||
).end(end_time=dequeued_at_ns)
|
||||
record_write_queue_wait(self.name, dequeued_at_ns - task.enqueued_at_ns)
|
||||
if conn_exception is not None:
|
||||
# fn never runs in this branch, so there is nothing to
|
||||
# wrap in a db.write.execute span.
|
||||
exception = conn_exception
|
||||
elif task.isolated_connection:
|
||||
try:
|
||||
with tracer.start_as_current_span(
|
||||
DB_WRITE_EXECUTE, **write_span_kwargs
|
||||
) as span:
|
||||
span.set_attribute(
|
||||
ISOLATED_CONNECTION,
|
||||
task.isolated_connection,
|
||||
)
|
||||
span.set_attribute(TRANSACTION, task.transaction)
|
||||
isolated_connection = self.connect(write=True)
|
||||
try:
|
||||
result = task.fn(isolated_connection)
|
||||
finally:
|
||||
isolated_connection.close()
|
||||
try:
|
||||
self._all_file_connections.remove(
|
||||
isolated_connection
|
||||
)
|
||||
except ValueError:
|
||||
# Was probably a memory connection
|
||||
pass
|
||||
except Exception as e: # noqa: BLE001
|
||||
# Write thread must survive any task failure or the database wedges
|
||||
sys.stderr.write(f"{e}\n")
|
||||
sys.stderr.flush()
|
||||
exception = e
|
||||
else:
|
||||
try:
|
||||
with tracer.start_as_current_span(
|
||||
DB_WRITE_EXECUTE, **write_span_kwargs
|
||||
) as span:
|
||||
span.set_attribute(
|
||||
ISOLATED_CONNECTION,
|
||||
task.isolated_connection,
|
||||
)
|
||||
span.set_attribute(TRANSACTION, task.transaction)
|
||||
if task.transaction:
|
||||
with conn:
|
||||
conn.execute("BEGIN IMMEDIATE")
|
||||
result = task.fn(conn)
|
||||
else:
|
||||
result = task.fn(conn)
|
||||
except Exception as e: # noqa: BLE001
|
||||
sys.stderr.write(f"{e}\n")
|
||||
sys.stderr.flush()
|
||||
exception = e
|
||||
_deliver_write_result(task, result, exception)
|
||||
finally:
|
||||
if token is not None:
|
||||
otel_context_api.detach(token)
|
||||
try:
|
||||
if task.transaction:
|
||||
with conn:
|
||||
conn.execute("BEGIN IMMEDIATE")
|
||||
result = task.fn(conn)
|
||||
else:
|
||||
result = task.fn(conn)
|
||||
except Exception as e: # noqa: BLE001
|
||||
sys.stderr.write(f"{e}\n")
|
||||
sys.stderr.flush()
|
||||
exception = e
|
||||
_deliver_write_result(task, result, exception)
|
||||
|
||||
async def execute_fn(self, fn):
|
||||
"""Run `fn(conn)` on a read connection, traced as one database call.
|
||||
|
||||
The public entry point for callback-style reads - plugins and core
|
||||
both use it to run arbitrary Python against a connection. It is
|
||||
instrumented exactly like `execute()`: one `db.query` span (with
|
||||
`datasette.callback` in place of `db.query.text`, since there is no
|
||||
SQL string to record) and a `db.query.execute` child covering the
|
||||
time actually spent on the worker thread. `execute()` itself calls
|
||||
`_execute_fn()` directly, so a SQL read never gets a second span.
|
||||
"""
|
||||
self._check_not_closed()
|
||||
|
||||
def fn_in_execute_span(conn):
|
||||
# Created on the worker thread; parents to the db.query span via
|
||||
# the copy_context() propagation in _execute_fn(). The gap
|
||||
# between the two spans is time spent waiting for a free thread.
|
||||
with tracer.start_as_current_span(DB_QUERY_EXECUTE):
|
||||
return fn(conn)
|
||||
|
||||
with tracer.start_as_current_span(DB_QUERY, kind=DB_QUERY.kind) as span:
|
||||
span.set_attribute(DB_SYSTEM, "sqlite")
|
||||
span.set_attribute(DB_NAMESPACE, self.name)
|
||||
span.set_attribute(CALLBACK, callback_name(fn))
|
||||
# Default exception handling applies, unlike execute(): there is
|
||||
# no log_sql_errors=False probing caller and no expected-timeout
|
||||
# budget on this path, so a raised exception is an error.
|
||||
with record_operation_duration(self.name, "read"):
|
||||
return await self._execute_fn(fn_in_execute_span)
|
||||
|
||||
async def _execute_fn(self, fn):
|
||||
self._check_not_closed()
|
||||
if self.ds.executor is None:
|
||||
# non-threaded mode
|
||||
|
|
@ -749,29 +512,7 @@ class Database:
|
|||
|
||||
with self._pending_execute_futures_lock:
|
||||
self._check_not_closed()
|
||||
# A fresh copy_context() is required per submit (not one shared
|
||||
# copy reused across calls): concurrent execution of the same
|
||||
# Context raises "RuntimeError: cannot enter context ...
|
||||
# already entered". This propagates the caller's otel context
|
||||
# (e.g. the enclosing db.query span) onto the worker thread.
|
||||
#
|
||||
# copy_context() is not selective: it also carries Datasette's own
|
||||
# ContextVars - _skip_permission_checks and _permission_check_cache
|
||||
# (datasette/permissions.py), _in_datasette_client (app.py) and,
|
||||
# until the hand-rolled tracer goes, trace_task_id (tracer.py) -
|
||||
# into worker threads, where they previously took their defaults.
|
||||
# That is safe, for two reasons. Nothing reads them on a worker
|
||||
# thread: the permission code that reads the first two is async and
|
||||
# only ever runs on the event loop. And Context.run() restores the
|
||||
# thread's previous context when the callable returns, so a value
|
||||
# cannot outlive the submit that carried it and reach the next task
|
||||
# on this shared pool - "skip permission checks" in particular can
|
||||
# never bleed from one request into another's query. Where a value
|
||||
# would be read - a plugin calling datasette.in_client() or trace()
|
||||
# from inside an execute_fn callable - seeing the submitting
|
||||
# request's value is the more accurate answer, not a leak.
|
||||
ctx = contextvars.copy_context()
|
||||
future = self.ds.executor.submit(ctx.run, in_thread)
|
||||
future = self.ds.executor.submit(in_thread)
|
||||
self._pending_execute_futures.add(future)
|
||||
future.add_done_callback(self._remove_pending_execute_future)
|
||||
return await asyncio.wrap_future(future)
|
||||
|
|
@ -784,154 +525,48 @@ class Database:
|
|||
custom_time_limit=None,
|
||||
page_size=None,
|
||||
log_sql_errors=True,
|
||||
table=None,
|
||||
):
|
||||
"""Executes sql against db_name in a thread
|
||||
|
||||
`table`, if passed, is recorded as the `db.collection.name` span
|
||||
attribute. It exists for callers that already know which table the
|
||||
query targets - the table and row views - and is never derived from
|
||||
`sql` itself: deriving it would be a parse, and on an instance where
|
||||
anyone can create a table the resulting value set has no ceiling.
|
||||
"""
|
||||
"""Executes sql against db_name in a thread"""
|
||||
self._check_not_closed()
|
||||
page_size = page_size or self.ds.page_size
|
||||
time_limit_ms = self.ds.sql_time_limit_ms
|
||||
# A caller that hands in a budget shorter than the instance-wide
|
||||
# sql_time_limit_ms is saying "this may not finish, and that is an
|
||||
# answer I can use" - and every such caller in core does treat the
|
||||
# timeout as normal: table_counts() stores None per table, facet
|
||||
# suggestion moves on to the next column, autocomplete falls back to a
|
||||
# prefix query. Those timeouts are therefore not span errors. Without
|
||||
# this, the homepage alone emits one red span per table (it counts
|
||||
# every table under a 10ms budget) on every single hit.
|
||||
#
|
||||
# A query that runs out the instance-wide limit is a different event -
|
||||
# nobody asked for a short budget, so it stays an error.
|
||||
timeout_expected = bool(custom_time_limit) and custom_time_limit < time_limit_ms
|
||||
if timeout_expected:
|
||||
time_limit_ms = custom_time_limit
|
||||
|
||||
def sql_operation_in_thread(conn):
|
||||
# This span is created inside the worker thread. Its parent is
|
||||
# resolved from the ambient otel context, which was propagated
|
||||
# onto this thread via copy_context() at the executor.submit()
|
||||
# boundary in _execute_fn() (or run_in_executor() for immutable
|
||||
# databases) - so it parents correctly to the enclosing
|
||||
# db.query span despite running on a different thread.
|
||||
#
|
||||
# Exception handling is explicit rather than left to the context
|
||||
# manager's flags, which apply to every exception type alike. This
|
||||
# span needs to tell two apart: an expected timeout is never an
|
||||
# error, while a genuine SQL failure is one unless the caller
|
||||
# passed log_sql_errors=False, meaning it was probing and treats
|
||||
# failure as an expected answer. Without the latter, facet
|
||||
# suggestion marks two spans per text column as failed on every
|
||||
# table page; without the former, so does every homepage hit.
|
||||
with tracer.start_as_current_span(
|
||||
DB_QUERY_EXECUTE,
|
||||
record_exception=False,
|
||||
set_status_on_exception=False,
|
||||
) as execute_span:
|
||||
time_limit_ms = self.ds.sql_time_limit_ms
|
||||
if custom_time_limit and custom_time_limit < time_limit_ms:
|
||||
time_limit_ms = custom_time_limit
|
||||
|
||||
with sqlite_timelimit(conn, time_limit_ms):
|
||||
try:
|
||||
with sqlite_timelimit(conn, time_limit_ms):
|
||||
try:
|
||||
cursor = conn.cursor()
|
||||
cursor.execute(sql, params if params is not None else {})
|
||||
max_returned_rows = self.ds.max_returned_rows
|
||||
if max_returned_rows == page_size:
|
||||
max_returned_rows += 1
|
||||
if max_returned_rows and truncate:
|
||||
rows = cursor.fetchmany(max_returned_rows + 1)
|
||||
truncated = len(rows) > max_returned_rows
|
||||
rows = rows[:max_returned_rows]
|
||||
else:
|
||||
rows = cursor.fetchall()
|
||||
truncated = False
|
||||
except (sqlite3.OperationalError, sqlite3.DatabaseError) as e:
|
||||
if e.args == ("interrupted",):
|
||||
raise QueryInterrupted(e, sql, params)
|
||||
if log_sql_errors:
|
||||
sys.stderr.write(
|
||||
f"ERROR: conn={conn}, sql = {sql!r}, params = {params}: {e}\n"
|
||||
)
|
||||
sys.stderr.flush()
|
||||
raise
|
||||
except QueryInterrupted as e:
|
||||
if not timeout_expected:
|
||||
execute_span.record_exception(e)
|
||||
execute_span.set_status(Status(StatusCode.ERROR, str(e)))
|
||||
raise
|
||||
except Exception as e:
|
||||
if log_sql_errors:
|
||||
execute_span.record_exception(e)
|
||||
execute_span.set_status(Status(StatusCode.ERROR, str(e)))
|
||||
raise
|
||||
|
||||
if truncate:
|
||||
return Results(rows, truncated, cursor.description)
|
||||
|
||||
else:
|
||||
return Results(rows, False, cursor.description)
|
||||
|
||||
# SIM117 wants these two context managers merged. They are kept nested
|
||||
# deliberately: the hand-rolled tracer's wrapper is on its way out, and
|
||||
# nesting makes removing it a single-line deletion.
|
||||
with trace( # noqa: SIM117
|
||||
"sql", database=self.name, sql=sql.strip(), params=params
|
||||
):
|
||||
# Exception handling is explicit rather than left to the context
|
||||
# manager's defaults, so that callers passing log_sql_errors=False
|
||||
# can be honoured - see the comment on the generic handler below.
|
||||
with tracer.start_as_current_span(
|
||||
DB_QUERY,
|
||||
kind=DB_QUERY.kind,
|
||||
record_exception=False,
|
||||
set_status_on_exception=False,
|
||||
) as span:
|
||||
span.set_attribute(DB_SYSTEM, "sqlite")
|
||||
span.set_attribute(DB_NAMESPACE, self.name)
|
||||
span.set_attribute(DB_QUERY_TEXT, sql_attribute(sql))
|
||||
span.set_attribute(TIME_LIMIT_MS, time_limit_ms)
|
||||
operation_name = sql_operation_name(sql)
|
||||
if operation_name:
|
||||
span.set_attribute(DB_OPERATION_NAME, operation_name)
|
||||
if table:
|
||||
span.set_attribute(DB_COLLECTION_NAME, table)
|
||||
if params:
|
||||
span.set_attribute(PARAM_COUNT, len(params))
|
||||
try:
|
||||
with record_operation_duration(self.name, "read"):
|
||||
results = await self._execute_fn(sql_operation_in_thread)
|
||||
except QueryInterrupted as e:
|
||||
# datasette.interrupted is set either way - it is the
|
||||
# signal worth having. Only the ERROR status is
|
||||
# conditional; see the timeout_expected comment above.
|
||||
span.set_attribute(INTERRUPTED, True)
|
||||
if not timeout_expected:
|
||||
span.set_status(Status(StatusCode.ERROR, str(e)))
|
||||
span.record_exception(e)
|
||||
# Expected timeouts (a caller that opted into a shorter
|
||||
# budget, like facet suggestion) are not counted - see
|
||||
# the M_QUERIES_INTERRUPTED registry entry for why.
|
||||
record_query_interrupted(self.name)
|
||||
raise
|
||||
except Exception as e:
|
||||
# log_sql_errors=False means the caller is probing and
|
||||
# treats failure as an expected answer, not an error.
|
||||
# Facet suggestion is the big one: it runs json_type()
|
||||
# against every column precisely to find out which ones
|
||||
# raise, so a table with N text columns would otherwise
|
||||
# mark N queries per page as failed - burying real errors
|
||||
# and setting off any alerting based on span status.
|
||||
if log_sql_errors:
|
||||
span.record_exception(e)
|
||||
span.set_status(Status(StatusCode.ERROR, str(e)))
|
||||
cursor = conn.cursor()
|
||||
cursor.execute(sql, params if params is not None else {})
|
||||
max_returned_rows = self.ds.max_returned_rows
|
||||
if max_returned_rows == page_size:
|
||||
max_returned_rows += 1
|
||||
if max_returned_rows and truncate:
|
||||
rows = cursor.fetchmany(max_returned_rows + 1)
|
||||
truncated = len(rows) > max_returned_rows
|
||||
rows = rows[:max_returned_rows]
|
||||
else:
|
||||
span.set_attribute(SQL_ERROR_SUPPRESSED, True)
|
||||
rows = cursor.fetchall()
|
||||
truncated = False
|
||||
except (sqlite3.OperationalError, sqlite3.DatabaseError) as e:
|
||||
if e.args == ("interrupted",):
|
||||
raise QueryInterrupted(e, sql, params)
|
||||
if log_sql_errors:
|
||||
sys.stderr.write(
|
||||
f"ERROR: conn={conn}, sql = {sql!r}, params = {params}: {e}\n"
|
||||
)
|
||||
sys.stderr.flush()
|
||||
raise
|
||||
span.set_attribute(TRUNCATED, results.truncated)
|
||||
span.set_attribute(ROWS_RETURNED, len(results.rows))
|
||||
|
||||
if truncate:
|
||||
return Results(rows, truncated, cursor.description)
|
||||
|
||||
else:
|
||||
return Results(rows, False, cursor.description)
|
||||
|
||||
with trace("sql", database=self.name, sql=sql.strip(), params=params):
|
||||
results = await self.execute_fn(sql_operation_in_thread)
|
||||
return results
|
||||
|
||||
@property
|
||||
|
|
@ -1022,34 +657,17 @@ class Database:
|
|||
)
|
||||
return [r[0] for r in results.rows]
|
||||
|
||||
# These callbacks are named functions rather than lambdas so that their
|
||||
# db.query spans carry a greppable datasette.callback - exactly the
|
||||
# guidance the plugin telemetry docs give, applied to core's own
|
||||
# highest-frequency introspection calls.
|
||||
|
||||
async def table_columns(self, table):
|
||||
def _table_columns(conn):
|
||||
return table_columns(conn, table)
|
||||
|
||||
return await self.execute_fn(_table_columns)
|
||||
return await self.execute_fn(lambda conn: table_columns(conn, table))
|
||||
|
||||
async def table_column_details(self, table):
|
||||
def _table_column_details(conn):
|
||||
return table_column_details(conn, table)
|
||||
|
||||
return await self.execute_fn(_table_column_details)
|
||||
return await self.execute_fn(lambda conn: table_column_details(conn, table))
|
||||
|
||||
async def primary_keys(self, table):
|
||||
def _primary_keys(conn):
|
||||
return detect_primary_keys(conn, table)
|
||||
|
||||
return await self.execute_fn(_primary_keys)
|
||||
return await self.execute_fn(lambda conn: detect_primary_keys(conn, table))
|
||||
|
||||
async def fts_table(self, table):
|
||||
def _fts_table(conn):
|
||||
return detect_fts(conn, table)
|
||||
|
||||
return await self.execute_fn(_fts_table)
|
||||
return await self.execute_fn(lambda conn: detect_fts(conn, table))
|
||||
|
||||
async def label_column_for_table(self, table):
|
||||
explicit_label_column = (await self.ds.table_config(self.name, table)).get(
|
||||
|
|
@ -1236,28 +854,16 @@ def _apply_write_wrapper(fn, wrapper_factory, track_event):
|
|||
|
||||
class WriteTask:
|
||||
__slots__ = (
|
||||
"block",
|
||||
"enqueued_at_ns",
|
||||
"fn",
|
||||
"isolated_connection",
|
||||
"loop",
|
||||
"otel_context",
|
||||
"reply_future",
|
||||
"task_id",
|
||||
"transaction",
|
||||
)
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
fn,
|
||||
task_id,
|
||||
loop,
|
||||
reply_future,
|
||||
isolated_connection,
|
||||
transaction,
|
||||
otel_context,
|
||||
enqueued_at_ns,
|
||||
block,
|
||||
self, fn, task_id, loop, reply_future, isolated_connection, transaction
|
||||
):
|
||||
self.fn = fn
|
||||
self.task_id = task_id
|
||||
|
|
@ -1265,14 +871,6 @@ class WriteTask:
|
|||
self.reply_future = reply_future
|
||||
self.isolated_connection = isolated_connection
|
||||
self.transaction = transaction
|
||||
self.otel_context = otel_context
|
||||
self.enqueued_at_ns = enqueued_at_ns
|
||||
# Whether the enqueueing caller awaits the reply future. Decides how
|
||||
# `_execute_writes` relates this task's spans to `otel_context`:
|
||||
# parent (block=True) or span-link target (block=False). See the
|
||||
# comment at the WriteTask construction site in
|
||||
# `_send_to_write_thread`.
|
||||
self.block = block
|
||||
|
||||
|
||||
def _deliver_write_result(task, result, exception):
|
||||
|
|
|
|||
|
|
@ -1,649 +0,0 @@
|
|||
"""
|
||||
OpenTelemetry integration for Datasette core.
|
||||
|
||||
Core depends on `opentelemetry-api` only. It never creates a
|
||||
`TracerProvider` or a `MeterProvider`, never configures an exporter, and
|
||||
never touches sampling - that is the responsibility of whoever is running
|
||||
Datasette (an `opentelemetry-instrument` agent, a future plugin, or a test
|
||||
harness).
|
||||
|
||||
With no provider installed every span produced here is a
|
||||
`NonRecordingSpan`. That is not free - a table page emits ~100 spans -
|
||||
but end-to-end page benchmarks put the overhead below their own
|
||||
run-to-run variation. Installing an SDK provider is what costs
|
||||
something measurable. Every metric instrument is likewise a no-op
|
||||
without a provider, and the observable-gauge callbacks are never
|
||||
invoked at all.
|
||||
"""
|
||||
|
||||
import contextvars
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
import weakref
|
||||
from contextlib import contextmanager
|
||||
|
||||
from opentelemetry import context as otel_context_api
|
||||
from opentelemetry import metrics as otel_metrics
|
||||
from opentelemetry import trace as otel_trace
|
||||
from opentelemetry.propagate import extract
|
||||
from opentelemetry.propagators.textmap import Getter
|
||||
from opentelemetry.trace import Link, SpanKind, Status, StatusCode, get_current_span
|
||||
|
||||
from .telemetry_registry import (
|
||||
DB_NAMESPACE,
|
||||
DB_SYSTEM,
|
||||
ERROR_TYPE,
|
||||
HTTP_REQUEST_METHOD,
|
||||
HTTP_RESPONSE_STATUS_CODE,
|
||||
INTERNAL_CLIENT,
|
||||
M_CONNECTIONS_OPEN,
|
||||
M_OPERATION_DURATION,
|
||||
M_QUERIES_INTERRUPTED,
|
||||
M_QUERIES_PENDING,
|
||||
M_THREADS_LIMIT,
|
||||
M_THREADS_QUEUE_DEPTH,
|
||||
M_WRITE_QUEUE_DEPTH,
|
||||
M_WRITE_QUEUE_WAIT,
|
||||
OPERATION,
|
||||
SERVER_ADDRESS,
|
||||
URL_PATH,
|
||||
URL_SCHEME,
|
||||
USER_AGENT_ORIGINAL,
|
||||
)
|
||||
from .version import __version__
|
||||
|
||||
# True while code is executing within a datasette.client request. Defined
|
||||
# here rather than in app.py (which owns its writers and the in_client()
|
||||
# accessor) so TelemetryMiddleware can read it without a circular import:
|
||||
# an in-process sub-request runs the full ASGI stack, so it emits a second,
|
||||
# nested SERVER span - datasette.internal_client marks those so kind-based
|
||||
# dashboards can filter the double-count out.
|
||||
_in_datasette_client = contextvars.ContextVar("in_datasette_client", default=False)
|
||||
|
||||
# The semantic-convention version whose spellings this instrumentation
|
||||
# actually emits. Deliberately NOT the latest release.
|
||||
#
|
||||
# A schema URL is a machine-readable claim: a consumer doing schema
|
||||
# translation replays the renames between the declared version and the one
|
||||
# it wants, so the claim has to name the version whose spellings are on the
|
||||
# wire. A wrong one makes translation wrong rather than merely uninformative.
|
||||
#
|
||||
# Datasette emits `db.system`, which was renamed to `db.system.name` in
|
||||
# semconv 1.30.0. Everything else it emits (`db.namespace`, `db.query.text`,
|
||||
# `db.operation.name`, `db.collection.name`) has been current since 1.26.0.
|
||||
# So 1.29.0 is the highest version at which every name emitted here is the
|
||||
# current spelling. Everything under `datasette.*` is Datasette's own and
|
||||
# outside semconv, so it is unaffected either way.
|
||||
#
|
||||
# Declaring 1.43.0 would be false about `db.system`, and would actively STOP
|
||||
# a consumer translating it forward, because it asserts the rename already
|
||||
# happened. Bump this deliberately, in the same commit as the attribute
|
||||
# renames it implies - it is a claim about the names, not decoration.
|
||||
SCHEMA_URL = "https://opentelemetry.io/schemas/1.29.0"
|
||||
|
||||
tracer = otel_trace.get_tracer("datasette", __version__, schema_url=SCHEMA_URL)
|
||||
meter = otel_metrics.get_meter("datasette", __version__, schema_url=SCHEMA_URL)
|
||||
|
||||
MAX_SQL_LENGTH = 2048
|
||||
|
||||
|
||||
def sql_attribute(sql: str) -> str:
|
||||
"Truncate SQL text so it is safe to attach to a span as an attribute."
|
||||
sql = sql.strip()
|
||||
if len(sql) <= MAX_SQL_LENGTH:
|
||||
return sql
|
||||
return sql[:MAX_SQL_LENGTH] + "…[truncated]"
|
||||
|
||||
|
||||
def callback_name(fn) -> str:
|
||||
"""
|
||||
The name recorded as `datasette.callback` for a callback-style call.
|
||||
|
||||
`functools.partial` objects (and other callables) have no `__qualname__`,
|
||||
so fall back to the type's name rather than fail the query over telemetry.
|
||||
"""
|
||||
return getattr(fn, "__qualname__", type(fn).__name__)
|
||||
|
||||
|
||||
def linked_root_span_kwargs(context=None):
|
||||
"""
|
||||
Keyword arguments that start a span as a root in its own trace, carrying
|
||||
a ``Link`` back to whatever span is current - the shape for work that a
|
||||
request *caused* without *containing*.
|
||||
|
||||
Use it when the causing span will end before the work does (a background
|
||||
task, a scheduled job, a ``block=False`` write): parenting there would
|
||||
draw a child outliving its closed parent, which renders badly in most
|
||||
trace UIs. The explicit empty ``Context()`` also stops the worker
|
||||
thread's ambient context from supplying an accidental parent.
|
||||
|
||||
Pass ``context`` to link to the span current in a *captured* context
|
||||
(e.g. one carried across a queue) rather than the caller's. If no valid
|
||||
span is current there is simply no link. The link carries no attributes:
|
||||
with only one kind of link, naming the relationship would add nothing.
|
||||
|
||||
Works with any tracer::
|
||||
|
||||
with my_tracer.start_as_current_span(
|
||||
"myplugin.job", **linked_root_span_kwargs()
|
||||
):
|
||||
...
|
||||
"""
|
||||
cause = get_current_span(context).get_span_context()
|
||||
links = [Link(cause)] if cause.is_valid else []
|
||||
return {"context": otel_context_api.Context(), "links": links}
|
||||
|
||||
|
||||
# db.operation.name is the leading keyword of a statement matched against a
|
||||
# fixed allowlist - deliberately not a parse.
|
||||
#
|
||||
# This runs against arbitrary user-supplied SQL (the `?sql=` query string,
|
||||
# canned queries, anything typed into the query editor), and the attribute
|
||||
# has to stay safe to use as a metric dimension. A metric series is keyed by
|
||||
# its attribute values, so echoing back an arbitrary first token would let
|
||||
# one visitor's typo mint a new, permanent series. The allowlist bounds that
|
||||
# at a fixed, small set regardless of what anyone sends.
|
||||
DB_OPERATION_ALLOWLIST = frozenset(
|
||||
{
|
||||
"SELECT",
|
||||
"INSERT",
|
||||
"UPDATE",
|
||||
"DELETE",
|
||||
"CREATE",
|
||||
"DROP",
|
||||
"ALTER",
|
||||
"PRAGMA",
|
||||
"EXPLAIN",
|
||||
"REPLACE",
|
||||
"VACUUM",
|
||||
"ANALYZE",
|
||||
"WITH",
|
||||
}
|
||||
)
|
||||
|
||||
_LEADING_KEYWORD = re.compile(r"^\s*([A-Za-z]+)")
|
||||
|
||||
|
||||
def sql_operation_name(sql: str) -> str | None:
|
||||
"""
|
||||
The statement's leading keyword, if it is one we recognise.
|
||||
|
||||
Returns None - never a guess - for anything not on the allowlist,
|
||||
including a statement that opens with a comment or with punctuation such
|
||||
as the "(" of a parenthesised SELECT.
|
||||
|
||||
Known limitation: a statement beginning with a CTE reports `WITH` rather
|
||||
than the operation inside it, and a substantial share of Datasette's own
|
||||
reads take that form. Extracting more than the leading keyword means
|
||||
handling comment stripping, parenthesised `(SELECT ...) UNION` and
|
||||
compound names like `CREATE TABLE` - each a special case a hand-rolled
|
||||
matcher would accrete and eventually get wrong. Omitting a name beats
|
||||
guessing at one.
|
||||
|
||||
Only safe to call with a single statement: `execute_write_script()` runs
|
||||
several separated by semicolons, and semantic conventions say
|
||||
`db.operation.name` "SHOULD NOT be extracted from db.query.text, when the
|
||||
database system supports query text with multiple operations in non-batch
|
||||
operations" - so that call site does not use this at all rather than
|
||||
reporting only the first statement's operation.
|
||||
"""
|
||||
match = _LEADING_KEYWORD.match(sql)
|
||||
if not match:
|
||||
return None
|
||||
keyword = match.group(1).upper()
|
||||
if keyword in DB_OPERATION_ALLOWLIST:
|
||||
return keyword
|
||||
return None
|
||||
|
||||
|
||||
# --- The HTTP request span ------------------------------------------------
|
||||
|
||||
|
||||
class _ScopeHeadersGetter(Getter):
|
||||
"""
|
||||
Read W3C trace context out of an ASGI scope's headers.
|
||||
|
||||
`scope["headers"]` is a list of `(bytes, bytes)` pairs, lowercased by the
|
||||
server per the ASGI spec - but `.lower()` is applied again here because
|
||||
that is a spec promise about servers, not something this process
|
||||
controls. Header bytes are latin-1 by RFC 9110.
|
||||
"""
|
||||
|
||||
def get(self, carrier, key):
|
||||
wanted = key.lower().encode("latin-1")
|
||||
values = [v.decode("latin-1") for k, v in carrier if k.lower() == wanted]
|
||||
return values or None
|
||||
|
||||
def keys(self, carrier):
|
||||
return [k.decode("latin-1") for k, _ in carrier]
|
||||
|
||||
|
||||
_HEADERS_GETTER = _ScopeHeadersGetter()
|
||||
|
||||
|
||||
# An unclamped method is an unbounded dimension a client controls: anyone can
|
||||
# send `FOO / HTTP/1.1`. Semantic conventions say map anything unrecognised to
|
||||
# `_OTHER`. These nine are the methods of RFC 9110 plus PATCH (RFC 5789).
|
||||
_KNOWN_METHODS = frozenset(
|
||||
{"GET", "HEAD", "POST", "PUT", "DELETE", "CONNECT", "OPTIONS", "TRACE", "PATCH"}
|
||||
)
|
||||
|
||||
|
||||
def clamp_http_method(method):
|
||||
"The request method if it is one we recognise, else ``_OTHER``."
|
||||
method = (method or "").upper()
|
||||
return method if method in _KNOWN_METHODS else "_OTHER"
|
||||
|
||||
|
||||
def _first_header(headers, name):
|
||||
"The first value of a header, decoded, or None."
|
||||
for key, value in headers:
|
||||
if key.lower() == name:
|
||||
return value.decode("latin-1")
|
||||
return None
|
||||
|
||||
|
||||
def _url_path(scope):
|
||||
"""
|
||||
The request path, with any query string removed.
|
||||
|
||||
`raw_path` is preferred because it is the bytes the client sent, before
|
||||
percent-decoding - Datasette routes on database and table names that can
|
||||
contain encoded slashes, which `scope["path"]` has already collapsed.
|
||||
|
||||
The split on "?" is not decoration. The ASGI spec's `raw_path` excludes
|
||||
the query string, but the name is read both ways in the wild - httpx's
|
||||
own `raw_path` includes the query - and Datasette's query strings carry
|
||||
user-supplied SQL, which core never records. A literal "?" cannot appear
|
||||
unencoded in a path, so the defensive split costs nothing when the server
|
||||
is well behaved.
|
||||
"""
|
||||
raw_path = scope.get("raw_path")
|
||||
if raw_path:
|
||||
if isinstance(raw_path, bytes):
|
||||
raw_path = raw_path.decode("latin-1")
|
||||
return raw_path.split("?", 1)[0]
|
||||
return scope.get("path", "")
|
||||
|
||||
|
||||
# The request span is handed to `DatasetteRouter.route_path` through the ASGI
|
||||
# scope rather than through `get_current_span()`, because by the time routing
|
||||
# happens the current span may well be something else: a plugin
|
||||
# `asgi_wrapper()` runs *inside* this middleware, and an instrumented one makes
|
||||
# its own span current for the whole request. Reading the current span there
|
||||
# would set `http.route` on that plugin's span - and rename it - while leaving
|
||||
# the actual request span without the one attribute a trace UI groups by. Not
|
||||
# hypothetical: an ordinary tracing plugin triggers it.
|
||||
#
|
||||
# Namespaced per the ASGI spec's rules for extension keys. Absent when the span
|
||||
# is not recording, which is exactly when the router should skip the work too.
|
||||
REQUEST_SPAN_SCOPE_KEY = "datasette.telemetry.request_span"
|
||||
|
||||
|
||||
def request_span(scope):
|
||||
"""
|
||||
The recording request span for an ASGI scope, or None.
|
||||
|
||||
Falls back to the current span so that a `DatasetteRouter` running under
|
||||
some other instrumentation - one that started a SERVER span but of course
|
||||
knows nothing about this scope key - still gets enriched.
|
||||
"""
|
||||
span = scope.get(REQUEST_SPAN_SCOPE_KEY)
|
||||
if span is None:
|
||||
span = otel_trace.get_current_span()
|
||||
# is_recording(), not `get_span_context().is_valid` - see the fast-path
|
||||
# comment in TelemetryMiddleware for why valid is not the same as recording.
|
||||
return span if span.is_recording() else None
|
||||
|
||||
|
||||
class TelemetryMiddleware:
|
||||
"""
|
||||
One `SpanKind.SERVER` span per HTTP request.
|
||||
|
||||
Mounted outermost in `Datasette.app()`, so every other span raised while
|
||||
serving a request - database queries, plugin middleware, startup work on
|
||||
a cold ASGI-hosted deployment - has somewhere to belong instead of
|
||||
becoming its own root trace.
|
||||
|
||||
Deliberately much smaller than `opentelemetry-instrumentation-asgi`,
|
||||
which needs several hundred lines of deferred-end machinery for
|
||||
applications that return before their body is sent. Datasette does not:
|
||||
`DatasetteRouter.route_path` awaits `response.asgi_send(send)`, and for a
|
||||
streaming CSV export `AsgiStream.asgi_send` runs the generator inline.
|
||||
All of it happens inside the single `await self.app(...)` below, so
|
||||
ending the span in a `finally` covers the response body too.
|
||||
"""
|
||||
|
||||
def __init__(self, app):
|
||||
self.app = app
|
||||
|
||||
async def __call__(self, scope, receive, send):
|
||||
# First, before anything else: `AsgiLifespan` is *inside* this
|
||||
# middleware, so lifespan startup and shutdown have to pass through
|
||||
# untouched or the server never starts. Same for websockets.
|
||||
if scope["type"] != "http":
|
||||
await self.app(scope, receive, send)
|
||||
return
|
||||
headers = scope.get("headers") or []
|
||||
# The *global* propagator, deliberately: it leaves the operator in
|
||||
# control with no Datasette-specific setting - OTEL_PROPAGATORS=none
|
||||
# disables extraction entirely, OTEL_PROPAGATORS=tracecontext drops
|
||||
# baggage - and core configuring propagation itself would be the same
|
||||
# mistake as core configuring sampling.
|
||||
context = extract(headers, getter=_HEADERS_GETTER)
|
||||
method = clamp_http_method(scope.get("method", ""))
|
||||
# The method, not the URL: a span name has to be low cardinality, and
|
||||
# the method is what is known out here at the edge, before any routing
|
||||
# has happened.
|
||||
with tracer.start_as_current_span(
|
||||
method, context=context, kind=SpanKind.SERVER
|
||||
) as span:
|
||||
if not span.is_recording():
|
||||
# No provider installed, or a sampler dropped this trace.
|
||||
# Everything below would be discarded, so skip building the
|
||||
# `send` wrapper and let a default install pay almost
|
||||
# nothing. Note this cannot be `get_span_context().is_valid`:
|
||||
# with no provider but an inbound `traceparent`, the API's
|
||||
# NoOpTracer returns a NonRecordingSpan carrying the *remote*
|
||||
# context, which is perfectly valid and still records nothing.
|
||||
await self.app(scope, receive, send)
|
||||
return
|
||||
span.set_attribute(HTTP_REQUEST_METHOD, method)
|
||||
span.set_attribute(URL_PATH, _url_path(scope))
|
||||
scheme = scope.get("scheme")
|
||||
if scheme:
|
||||
span.set_attribute(URL_SCHEME, scheme)
|
||||
host = _first_header(headers, b"host")
|
||||
if host:
|
||||
span.set_attribute(SERVER_ADDRESS, host)
|
||||
user_agent = _first_header(headers, b"user-agent")
|
||||
if user_agent:
|
||||
span.set_attribute(USER_AGENT_ORIGINAL, user_agent)
|
||||
if _in_datasette_client.get():
|
||||
span.set_attribute(INTERNAL_CLIENT, True)
|
||||
|
||||
# A copy, not a mutation: the scope belongs to the server, and
|
||||
# every other layer in Datasette extends it the same way.
|
||||
scope = dict(scope, **{REQUEST_SPAN_SCOPE_KEY: span})
|
||||
|
||||
# The status cannot be read off a Response object: `asgi_static`,
|
||||
# the favicon route, `AsgiStream` and `AsgiFileDownload` all call
|
||||
# `send` directly and never build one. Wrapping `send` is the only
|
||||
# thing that sees every response, including the 404 and 500
|
||||
# handlers.
|
||||
status_holder = {}
|
||||
|
||||
async def wrapped_send(message):
|
||||
if (
|
||||
message["type"] == "http.response.start"
|
||||
and "status" not in status_holder
|
||||
):
|
||||
status_holder["status"] = message["status"]
|
||||
await send(message)
|
||||
|
||||
escaped = False
|
||||
try:
|
||||
await self.app(scope, receive, wrapped_send)
|
||||
except BaseException as exception:
|
||||
# BaseException, not Exception: `route_path` turns almost
|
||||
# everything into a 500 itself, but `asyncio.CancelledError`
|
||||
# on client disconnect is a BaseException its `except
|
||||
# Exception` deliberately does not catch.
|
||||
escaped = True
|
||||
span.set_attribute(ERROR_TYPE, type(exception).__name__)
|
||||
span.set_status(Status(StatusCode.ERROR, str(exception)))
|
||||
raise
|
||||
finally:
|
||||
status = status_holder.get("status")
|
||||
if status is not None:
|
||||
span.set_attribute(HTTP_RESPONSE_STATUS_CODE, status)
|
||||
# 4xx is NOT an error for a SERVER span per semantic
|
||||
# conventions - the client made the mistake, not us.
|
||||
#
|
||||
# `not escaped` because this block still runs when an
|
||||
# exception is on its way out, and a response can have
|
||||
# started before it: the exception's class name is more
|
||||
# use than the string "500", so it wins.
|
||||
if status >= 500 and not escaped:
|
||||
span.set_status(Status(StatusCode.ERROR))
|
||||
span.set_attribute(ERROR_TYPE, str(status))
|
||||
|
||||
|
||||
# --- Metrics --------------------------------------------------------------
|
||||
#
|
||||
# Two shapes. Observable gauges - a callback the SDK invokes on its own
|
||||
# collection cycle, so a no-provider install never runs them - answer level
|
||||
# questions no span can, like "am I saturating my SQL threads right now".
|
||||
# Synchronous histograms/counters are recorded inline on the query path and
|
||||
# survive trace sampling: 1% of traces still means 100% of the latency
|
||||
# distribution. Why each metric exists is documented on its registry entry.
|
||||
#
|
||||
# Note a real difference from tracing: `_ProxyMeter` and its instruments
|
||||
# forward to a provider installed *after* they were created, whereas
|
||||
# `ProxyTracer` permanently caches the concrete tracer it first resolves. So
|
||||
# module-level instruments here are safe, and tests do not need a provider
|
||||
# installed before this module is imported.
|
||||
|
||||
|
||||
def _duration_attributes(database_name, operation):
|
||||
return {
|
||||
DB_SYSTEM: "sqlite",
|
||||
DB_NAMESPACE: database_name,
|
||||
OPERATION: operation,
|
||||
}
|
||||
|
||||
|
||||
# Each instrument passes the SDK a short plain-text description; the registry
|
||||
# entry for the same metric carries a longer RST one for the generated docs
|
||||
# (it can use `:ref:` roles, which an exported description string cannot).
|
||||
|
||||
sql_operation_duration = meter.create_histogram(
|
||||
M_OPERATION_DURATION,
|
||||
unit=M_OPERATION_DURATION.unit,
|
||||
description="Duration of a SQL operation issued by Datasette",
|
||||
explicit_bucket_boundaries_advisory=M_OPERATION_DURATION.buckets,
|
||||
)
|
||||
|
||||
write_queue_wait = meter.create_histogram(
|
||||
M_WRITE_QUEUE_WAIT,
|
||||
unit=M_WRITE_QUEUE_WAIT.unit,
|
||||
description=(
|
||||
"Time a write spent queued behind the single write thread for its database"
|
||||
),
|
||||
explicit_bucket_boundaries_advisory=M_WRITE_QUEUE_WAIT.buckets,
|
||||
)
|
||||
|
||||
queries_interrupted = meter.create_counter(
|
||||
M_QUERIES_INTERRUPTED,
|
||||
unit=M_QUERIES_INTERRUPTED.unit,
|
||||
description=(
|
||||
"Queries cancelled for exceeding sql_time_limit_ms. Not derivable from "
|
||||
"spans under sampling, and the signal that a time limit is too tight"
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@contextmanager
|
||||
def record_operation_duration(database_name, operation):
|
||||
"""
|
||||
Record `db.client.operation.duration` for one SQL operation.
|
||||
|
||||
`error.type` is set from the exception class on failure, per semconv, so a
|
||||
latency distribution can be split by success and failure. For a
|
||||
`block=False` write this measures the enqueue, not the write - the same
|
||||
caveat that applies to the surrounding span.
|
||||
"""
|
||||
attributes = _duration_attributes(database_name, operation)
|
||||
started = time.perf_counter()
|
||||
try:
|
||||
yield
|
||||
except BaseException as exception:
|
||||
attributes[ERROR_TYPE] = type(exception).__qualname__
|
||||
raise
|
||||
finally:
|
||||
sql_operation_duration.record(time.perf_counter() - started, attributes)
|
||||
|
||||
|
||||
def record_write_queue_wait(database_name, waited_ns):
|
||||
write_queue_wait.record(waited_ns / 1e9, {DB_NAMESPACE: database_name})
|
||||
|
||||
|
||||
def record_query_interrupted(database_name):
|
||||
queries_interrupted.add(1, {DB_NAMESPACE: database_name})
|
||||
|
||||
|
||||
# Live Datasette instances, weakly held so that instrumenting an instance
|
||||
# never keeps it alive. Guarded by a lock because the gauge callbacks run on
|
||||
# the SDK's collection thread while the event loop may be building or closing
|
||||
# a Datasette.
|
||||
#
|
||||
# Known limitation: the pool gauges below carry no attribute identifying which
|
||||
# Datasette produced them, so if a single process runs more than one instance
|
||||
# their observations collide and last-one-wins. Production runs one instance
|
||||
# per process; adding an instance id to make the test suite's hundreds of
|
||||
# instances distinguishable would mean unbounded attribute cardinality in
|
||||
# exchange for fixing a case that does not occur in production.
|
||||
_live_datasettes = weakref.WeakSet()
|
||||
_live_datasettes_lock = threading.Lock()
|
||||
|
||||
|
||||
def register_datasette(ds):
|
||||
"Start reporting pool/queue gauges for this Datasette instance."
|
||||
with _live_datasettes_lock:
|
||||
_live_datasettes.add(ds)
|
||||
|
||||
|
||||
def unregister_datasette(ds):
|
||||
"Stop reporting gauges for an instance that has been closed."
|
||||
with _live_datasettes_lock:
|
||||
_live_datasettes.discard(ds)
|
||||
|
||||
|
||||
def _live_instances():
|
||||
with _live_datasettes_lock:
|
||||
return list(_live_datasettes)
|
||||
|
||||
|
||||
def _databases_of(ds):
|
||||
"""
|
||||
Every Database attached to an instance, including the internal database.
|
||||
|
||||
The internal database is deliberately included: permission checks run SQL
|
||||
against it on essentially every request, so its queue depth and connection
|
||||
count are as operationally interesting as any user database's.
|
||||
"""
|
||||
databases = list(ds.databases.values())
|
||||
internal = getattr(ds, "_internal_database", None)
|
||||
if internal is not None:
|
||||
databases.append(internal)
|
||||
return databases
|
||||
|
||||
|
||||
# Each callback is a plain generator function so it can be unit-tested
|
||||
# directly, without standing up an SDK provider and a metric reader.
|
||||
|
||||
|
||||
def observe_sql_thread_limit(options=None):
|
||||
"Size of the shared read-query thread pool (the num_sql_threads setting)."
|
||||
for ds in _live_instances():
|
||||
if ds.executor is None:
|
||||
# num_sql_threads=0 - queries run on the event loop, no pool.
|
||||
continue
|
||||
yield otel_metrics.Observation(ds.setting("num_sql_threads"), {})
|
||||
|
||||
|
||||
def observe_sql_thread_queue_depth(options=None):
|
||||
"""
|
||||
Read queries waiting for a free thread in the shared pool.
|
||||
|
||||
This is the saturation signal: sustained above zero means requests are
|
||||
queueing on num_sql_threads. `_work_queue` is a private attribute of
|
||||
ThreadPoolExecutor, so its absence is tolerated rather than fatal - a
|
||||
missing gauge is much better than a crashed collection cycle.
|
||||
"""
|
||||
for ds in _live_instances():
|
||||
if ds.executor is None:
|
||||
continue
|
||||
work_queue = getattr(ds.executor, "_work_queue", None)
|
||||
if work_queue is None:
|
||||
continue
|
||||
yield otel_metrics.Observation(work_queue.qsize(), {})
|
||||
|
||||
|
||||
def observe_pending_queries(options=None):
|
||||
"""
|
||||
Read queries submitted to the pool and not yet finished, per database.
|
||||
|
||||
Summed across databases and compared against the thread limit, this is the
|
||||
utilisation half of the saturation picture. `len()` is deliberately taken
|
||||
without `_pending_execute_futures_lock`: it is atomic, and taking a lock
|
||||
held on the request path from the collection thread would let telemetry
|
||||
add latency to queries.
|
||||
"""
|
||||
for ds in _live_instances():
|
||||
for db in _databases_of(ds):
|
||||
yield otel_metrics.Observation(
|
||||
len(db._pending_execute_futures), {DB_NAMESPACE: db.name}
|
||||
)
|
||||
|
||||
|
||||
def observe_write_queue_depth(options=None):
|
||||
"""
|
||||
Writes queued behind the single write thread, per database.
|
||||
|
||||
Every database serialises its writes through one thread, so this is
|
||||
unbounded backpressure that no amount of num_sql_threads will relieve.
|
||||
"""
|
||||
for ds in _live_instances():
|
||||
for db in _databases_of(ds):
|
||||
write_queue = db._write_queue
|
||||
if write_queue is None:
|
||||
# No write has ever been queued for this database.
|
||||
continue
|
||||
yield otel_metrics.Observation(write_queue.qsize(), {DB_NAMESPACE: db.name})
|
||||
|
||||
|
||||
def observe_open_connections(options=None):
|
||||
"Open SQLite file connections tracked for closing, per database."
|
||||
for ds in _live_instances():
|
||||
for db in _databases_of(ds):
|
||||
yield otel_metrics.Observation(
|
||||
len(db._all_file_connections), {DB_NAMESPACE: db.name}
|
||||
)
|
||||
|
||||
|
||||
sql_thread_limit_gauge = meter.create_observable_gauge(
|
||||
M_THREADS_LIMIT,
|
||||
callbacks=[observe_sql_thread_limit],
|
||||
unit=M_THREADS_LIMIT.unit,
|
||||
description="Maximum concurrent read queries (the num_sql_threads setting)",
|
||||
)
|
||||
|
||||
sql_thread_queue_depth_gauge = meter.create_observable_gauge(
|
||||
M_THREADS_QUEUE_DEPTH,
|
||||
callbacks=[observe_sql_thread_queue_depth],
|
||||
unit=M_THREADS_QUEUE_DEPTH.unit,
|
||||
description="Read queries waiting for a free thread in the shared SQL pool",
|
||||
)
|
||||
|
||||
pending_queries_gauge = meter.create_observable_gauge(
|
||||
M_QUERIES_PENDING,
|
||||
callbacks=[observe_pending_queries],
|
||||
unit=M_QUERIES_PENDING.unit,
|
||||
description="Read queries submitted to the pool and not yet complete",
|
||||
)
|
||||
|
||||
write_queue_depth_gauge = meter.create_observable_gauge(
|
||||
M_WRITE_QUEUE_DEPTH,
|
||||
callbacks=[observe_write_queue_depth],
|
||||
unit=M_WRITE_QUEUE_DEPTH.unit,
|
||||
description="Writes queued behind a database's single write thread",
|
||||
)
|
||||
|
||||
open_connections_gauge = meter.create_observable_gauge(
|
||||
M_CONNECTIONS_OPEN,
|
||||
callbacks=[observe_open_connections],
|
||||
unit=M_CONNECTIONS_OPEN.unit,
|
||||
description="Open SQLite file connections tracked for closing",
|
||||
)
|
||||
|
|
@ -1,624 +0,0 @@
|
|||
"""
|
||||
The single source of truth for every span and span attribute that Datasette
|
||||
core emits.
|
||||
|
||||
Three things read this module, which is the point of it existing:
|
||||
|
||||
1. **The instrumentation itself.** `Attribute` and `SpanName` subclass `str`,
|
||||
so a registry entry *is* the string OpenTelemetry wants. Call sites pass
|
||||
`DB_NAMESPACE` where they used to pass `"db.namespace"` - no wrapper API
|
||||
over the OTel calls, no parallel structure to keep in step, and a typo is
|
||||
now an `ImportError` instead of a silently misnamed attribute.
|
||||
|
||||
2. **The documentation.** `docs/telemetry_doc.py` renders the span reference
|
||||
in `docs/internals.rst` from these definitions using cog, and
|
||||
`cog --check` runs in CI - so the docs cannot drift from the code.
|
||||
|
||||
3. **A conformance test.** `tests/test_telemetry_registry.py` makes real
|
||||
requests, collects every span and attribute actually emitted, and compares
|
||||
both directions: emitted-but-unregistered catches instrumentation added
|
||||
without documentation, registered-but-never-emitted catches documentation
|
||||
describing something that no longer exists. Neither the type system nor
|
||||
the generated docs can catch that second case.
|
||||
"""
|
||||
|
||||
from opentelemetry.trace import SpanKind
|
||||
|
||||
|
||||
class Attribute(str):
|
||||
"""
|
||||
A span attribute key, carrying its own documentation.
|
||||
|
||||
Subclasses `str` so it can be handed straight to `set_attribute()`.
|
||||
|
||||
Part of Datasette's public plugin API - plugins declare their own
|
||||
telemetry registries with these classes. See the "Telemetry for plugin
|
||||
authors" documentation.
|
||||
"""
|
||||
|
||||
__slots__ = ("description", "optional", "values")
|
||||
|
||||
def __new__(cls, name, description, optional=False, values=None):
|
||||
self = super().__new__(cls, name)
|
||||
self.description = description
|
||||
self.optional = optional
|
||||
# A closed enum vocabulary for the attribute's values, or None for
|
||||
# an open value set. Declaring one does two things: the conformance
|
||||
# helpers assert every emitted value is a member, and it marks the
|
||||
# attribute as bounded - safe to use as a metric dimension, where an
|
||||
# open value set would be a cardinality hazard.
|
||||
self.values = frozenset(values) if values is not None else None
|
||||
return self
|
||||
|
||||
def __repr__(self):
|
||||
return f"Attribute({str(self)!r})"
|
||||
|
||||
|
||||
class SpanName(str):
|
||||
"""A span name, carrying its documentation and the attributes it may set.
|
||||
|
||||
Part of Datasette's public plugin API, like `Attribute`.
|
||||
"""
|
||||
|
||||
__slots__ = ("attributes", "description", "dynamic", "kind", "prefix")
|
||||
|
||||
def __new__(
|
||||
cls,
|
||||
name,
|
||||
description,
|
||||
attributes=(),
|
||||
prefix=False,
|
||||
dynamic=False,
|
||||
kind=SpanKind.INTERNAL,
|
||||
):
|
||||
self = super().__new__(cls, name)
|
||||
self.description = description
|
||||
self.attributes = tuple(attributes)
|
||||
# True for a span family whose emitted names carry a variable suffix
|
||||
# after a fixed prefix - e.g. a plugin's `chat {model}` registered as
|
||||
# SpanName("chat ", ..., prefix=True) - so `span_for()` matches by
|
||||
# prefix rather than equality. Core registers none itself; the flag
|
||||
# exists for plugin registries.
|
||||
self.prefix = prefix
|
||||
# True when the emitted name is composed at runtime and shares no
|
||||
# fixed prefix with the registry entry - the HTTP request span, whose
|
||||
# name is the request method followed by the matched route. There is
|
||||
# no substring of the entry that could be matched against the wire, so
|
||||
# `span_for()` resolves these by span kind instead, and the entry's own
|
||||
# string is a template written for a human reading the generated
|
||||
# reference.
|
||||
self.dynamic = dynamic
|
||||
# SpanKind.INTERNAL by default - every span Datasette emits describes
|
||||
# its own internal work. db.query is the one exception: it is a real
|
||||
# database call, so semantic conventions (and trace UIs, which key
|
||||
# their database styling off this) expect SpanKind.CLIENT.
|
||||
self.kind = kind
|
||||
return self
|
||||
|
||||
def __repr__(self):
|
||||
return f"SpanName({str(self)!r})"
|
||||
|
||||
|
||||
class MetricName(str):
|
||||
"A metric name, carrying its instrument kind, unit and attributes."
|
||||
|
||||
__slots__ = ("attributes", "buckets", "description", "kind", "unit")
|
||||
|
||||
def __new__(cls, name, kind, unit, description, attributes=(), buckets=None):
|
||||
self = super().__new__(cls, name)
|
||||
self.kind = kind
|
||||
self.unit = unit
|
||||
self.description = description
|
||||
self.attributes = tuple(attributes)
|
||||
# Explicit histogram bucket boundaries, for histograms only. Passed to
|
||||
# create_histogram() as explicit_bucket_boundaries_advisory and
|
||||
# published in the generated docs, since an operator writing a
|
||||
# histogram_quantile() query needs to know them.
|
||||
self.buckets = tuple(buckets) if buckets is not None else None
|
||||
return self
|
||||
|
||||
def __repr__(self):
|
||||
return f"MetricName({str(self)!r})"
|
||||
|
||||
|
||||
COUNTER = "Counter"
|
||||
UPDOWN_COUNTER = "UpDownCounter"
|
||||
HISTOGRAM = "Histogram"
|
||||
GAUGE = "Observable gauge"
|
||||
|
||||
|
||||
# --- Attributes -----------------------------------------------------------
|
||||
#
|
||||
# Shared attributes are defined once and referenced by every span that sets
|
||||
# them, so "which spans carry db.namespace?" is answerable by grep.
|
||||
|
||||
HTTP_REQUEST_METHOD = Attribute(
|
||||
"http.request.method",
|
||||
"The HTTP method, clamped to the nine methods RFC 9110 and RFC 5789 "
|
||||
"define. Anything else is reported as ``_OTHER``: the method is a "
|
||||
"client-controlled string, so echoing it back unbounded would be a "
|
||||
"cardinality hazard.",
|
||||
)
|
||||
HTTP_RESPONSE_STATUS_CODE = Attribute(
|
||||
"http.response.status_code",
|
||||
"The status of the response, read from the ASGI ``http.response.start`` "
|
||||
"message rather than from a :ref:`internals_response` object - several "
|
||||
"views, including static files, file downloads and streaming CSV, send "
|
||||
"that message themselves and never build one. Omitted if the connection "
|
||||
"closed before anything was sent.",
|
||||
optional=True,
|
||||
)
|
||||
HTTP_ROUTE = Attribute(
|
||||
"http.route",
|
||||
"The route the request matched, as the compiled regular expression "
|
||||
"pattern Datasette routes with - for example "
|
||||
"``/(?P<database>[^\\/\\.]+)/(?P<table>[^\\/\\.]+)(\\.(?P<format>\\w+))?$`` "
|
||||
"for a table page. It is deliberately the pattern rather than a prettified "
|
||||
"``/{database}/{table}`` template: the route table is fixed when the app "
|
||||
"is built, so the pattern is exact, bounded and needs no parsing, whereas "
|
||||
"the transform into something prettier accretes edge cases. Unlike "
|
||||
"``url.path`` this is low cardinality, so it is the attribute to group by. "
|
||||
"Omitted when no route matched - a 404 - which is also when the span name "
|
||||
"falls back to the bare method.",
|
||||
optional=True,
|
||||
)
|
||||
URL_PATH = Attribute(
|
||||
"url.path",
|
||||
"The path portion of the URL. The query string is deliberately **not** "
|
||||
"recorded, on this or any other span: Datasette puts user-supplied SQL in "
|
||||
"``?sql=`` and canned query parameters in the query string, so exporting "
|
||||
"it by default would export exactly the data the rest of this "
|
||||
"instrumentation is careful with.",
|
||||
)
|
||||
URL_SCHEME = Attribute("url.scheme", "``http`` or ``https``.")
|
||||
SERVER_ADDRESS = Attribute(
|
||||
"server.address",
|
||||
"The ``Host`` header, verbatim - including any ``:port`` suffix, a "
|
||||
"deliberate deviation from semantic conventions' ``server.address`` / "
|
||||
"``server.port`` split. Client-controlled, so treat it as untrusted input "
|
||||
"rather than as the identity of the server.",
|
||||
optional=True,
|
||||
)
|
||||
USER_AGENT_ORIGINAL = Attribute(
|
||||
"user_agent.original",
|
||||
"The ``User-Agent`` header, verbatim. Omitted if the client sent none.",
|
||||
optional=True,
|
||||
)
|
||||
INTERNAL_CLIENT = Attribute(
|
||||
"datasette.internal_client",
|
||||
"``True`` when the request was made in-process through "
|
||||
"``datasette.client`` rather than arriving over the network. Such a "
|
||||
"sub-request runs the full ASGI stack, so it emits its own nested "
|
||||
"``SERVER`` span inside the outer request's - filter on this attribute "
|
||||
"to keep kind-based dashboards from double-counting requests. Omitted "
|
||||
"for real inbound requests.",
|
||||
optional=True,
|
||||
)
|
||||
ERROR_TYPE = Attribute(
|
||||
"error.type",
|
||||
"Set when the request failed: the exception class name if one escaped the "
|
||||
"application, otherwise the status code as a string for a 5xx response. "
|
||||
"A 4xx does **not** set this and does not set an error status - per "
|
||||
"semantic conventions a client error is not a server span's failure.",
|
||||
optional=True,
|
||||
)
|
||||
|
||||
DB_SYSTEM = Attribute("db.system", "Always ``sqlite``.")
|
||||
DB_NAMESPACE = Attribute("db.namespace", "Name of the database being queried.")
|
||||
OPERATION = Attribute(
|
||||
"datasette.operation",
|
||||
"Whether the operation was a read or a write.",
|
||||
values={"read", "write"},
|
||||
)
|
||||
DB_QUERY_TEXT = Attribute(
|
||||
"db.query.text",
|
||||
"The SQL, truncated to 2048 characters. Never the parameter values. "
|
||||
"Absent for a callback-style call (``execute_fn()`` and friends), where "
|
||||
"there is no SQL string to record - ``datasette.callback`` is set "
|
||||
"instead.",
|
||||
optional=True,
|
||||
)
|
||||
CALLBACK = Attribute(
|
||||
"datasette.callback",
|
||||
"The qualified name of the Python callable passed to ``execute_fn()``, "
|
||||
"``execute_write_fn()`` or ``execute_isolated_fn()`` - for example "
|
||||
"``TableInsertView.post.<locals>.insert_or_upsert_rows``. Set instead of "
|
||||
"``db.query.text``, which does not exist for a callback: the SQL is "
|
||||
"whatever the function chooses to run. A lambda reports ``<lambda>``, "
|
||||
"which is why callers wanting a recognisable span should pass a named "
|
||||
"function. Bounded cardinality: the set of callables is fixed by the "
|
||||
"installed code, not by request input.",
|
||||
optional=True,
|
||||
)
|
||||
DB_OPERATION_NAME = Attribute(
|
||||
"db.operation.name",
|
||||
"The statement's leading keyword - ``SELECT``, ``INSERT``, ``CREATE``, and "
|
||||
"so on - matched against a small fixed allowlist. Omitted rather than set "
|
||||
"to an arbitrary value: the attribute must stay safe to use as a metric "
|
||||
"dimension, and echoing an unrecognised first token from user-supplied "
|
||||
"SQL would be an unbounded-cardinality hazard. Also omitted for "
|
||||
"``execute_write_script()``, which runs multiple statements - per "
|
||||
"semantic conventions, the operation name should not be extracted from "
|
||||
"query text that can contain more than one operation. Note that a "
|
||||
"statement beginning with a CTE reports ``WITH``, not the operation "
|
||||
"inside it - a substantial share of Datasette's own reads take that "
|
||||
"form. Resolving it further would mean parsing.",
|
||||
optional=True,
|
||||
)
|
||||
DB_COLLECTION_NAME = Attribute(
|
||||
"db.collection.name",
|
||||
"The primary table, set only where the view already knows it - the table "
|
||||
"and row pages. Omitted for arbitrary ``?sql=`` queries, where determining "
|
||||
"the table would mean parsing the query.",
|
||||
optional=True,
|
||||
)
|
||||
|
||||
PARAM_COUNT = Attribute(
|
||||
"datasette.param_count",
|
||||
"Number of bound parameters. Recorded instead of the values themselves.",
|
||||
optional=True,
|
||||
)
|
||||
PARAM_SETS = Attribute(
|
||||
"datasette.param_sets",
|
||||
"Number of parameter sets consumed by ``execute_write_many()``. Not a row "
|
||||
"count - ``executemany()`` returns no rows. The parameter values "
|
||||
"themselves are never recorded: that sequence can hold thousands of rows.",
|
||||
optional=True,
|
||||
)
|
||||
TIME_LIMIT_MS = Attribute(
|
||||
"datasette.time_limit_ms",
|
||||
"The :ref:`setting_sql_time_limit_ms` value this query ran under. Set on "
|
||||
"reads, which are the queries that time limit applies to.",
|
||||
optional=True,
|
||||
)
|
||||
ROWS_RETURNED = Attribute(
|
||||
"datasette.rows_returned",
|
||||
"Number of rows a read returned. Set on the read path only, and only when "
|
||||
"the read succeeded.",
|
||||
optional=True,
|
||||
)
|
||||
TRUNCATED = Attribute(
|
||||
"datasette.truncated",
|
||||
"True if the result was cut short by :ref:`setting_max_returned_rows`.",
|
||||
optional=True,
|
||||
)
|
||||
INTERRUPTED = Attribute(
|
||||
"datasette.interrupted",
|
||||
"True if the query was cancelled for exceeding the time limit. The span "
|
||||
"status is also set to ``ERROR``, unless the caller asked for a budget "
|
||||
"shorter than :ref:`setting_sql_time_limit_ms` - as table counts, facet "
|
||||
"suggestion and autocomplete all do - in which case running out of time "
|
||||
"is an expected answer rather than a failure and the status is left "
|
||||
"unset.",
|
||||
optional=True,
|
||||
)
|
||||
SQL_ERROR_SUPPRESSED = Attribute(
|
||||
"datasette.sql_error_suppressed",
|
||||
"True when the query failed but the caller passed ``log_sql_errors=False``, "
|
||||
"meaning it was probing and treats failure as an expected answer. Facet "
|
||||
"suggestion does this against every column.",
|
||||
optional=True,
|
||||
)
|
||||
EXECUTESCRIPT = Attribute(
|
||||
"datasette.executescript",
|
||||
"True for ``execute_write_script()``, which runs multiple statements.",
|
||||
optional=True,
|
||||
)
|
||||
EXECUTEMANY = Attribute(
|
||||
"datasette.executemany",
|
||||
"True for ``execute_write_many()``, which runs one statement against many "
|
||||
"parameter sets.",
|
||||
optional=True,
|
||||
)
|
||||
ISOLATED_CONNECTION = Attribute(
|
||||
"datasette.isolated_connection",
|
||||
"True if the write ran on its own connection rather than the shared write "
|
||||
"connection.",
|
||||
)
|
||||
TRANSACTION = Attribute(
|
||||
"datasette.transaction",
|
||||
"False for statements such as ``VACUUM`` that cannot run inside a transaction.",
|
||||
)
|
||||
|
||||
|
||||
# --- Spans ----------------------------------------------------------------
|
||||
|
||||
HTTP_REQUEST = SpanName(
|
||||
"{http.request.method} {http.route}",
|
||||
"One span per HTTP request, created by the outermost layer of the ASGI "
|
||||
"stack - so plugin ``asgi_wrapper()`` middleware, CSRF protection and "
|
||||
"every database span raised while serving the request all nest inside "
|
||||
"it. Without it each of those would be its own root trace. The span name "
|
||||
"is not a fixed string: it is the method followed by the matched route, "
|
||||
"and just the method for a request that matched no route. The span starts "
|
||||
"at the ASGI edge, before routing has happened, so it is named for the "
|
||||
"method there and renamed once the route is known. "
|
||||
"W3C ``traceparent`` and ``baggage`` headers are extracted using the "
|
||||
"global propagator, so a request arriving from an already-traced caller "
|
||||
"continues that trace; set ``OTEL_PROPAGATORS=none`` to turn that off, "
|
||||
"and strip those headers at your proxy if your instance is public.",
|
||||
(
|
||||
HTTP_REQUEST_METHOD,
|
||||
HTTP_ROUTE,
|
||||
URL_PATH,
|
||||
URL_SCHEME,
|
||||
SERVER_ADDRESS,
|
||||
USER_AGENT_ORIGINAL,
|
||||
HTTP_RESPONSE_STATUS_CODE,
|
||||
ERROR_TYPE,
|
||||
INTERNAL_CLIENT,
|
||||
),
|
||||
dynamic=True,
|
||||
kind=SpanKind.SERVER,
|
||||
)
|
||||
|
||||
DB_QUERY = SpanName(
|
||||
"db.query",
|
||||
"A SQL operation issued by Datasette, covering the full round trip "
|
||||
"including any time spent queued for a thread. Callback-style calls - "
|
||||
"``execute_fn()``, ``execute_write_fn()`` and ``execute_isolated_fn()`` - "
|
||||
"appear here too, distinguished by ``datasette.callback`` in place of "
|
||||
"``db.query.text``.",
|
||||
(
|
||||
DB_SYSTEM,
|
||||
DB_NAMESPACE,
|
||||
DB_QUERY_TEXT,
|
||||
CALLBACK,
|
||||
DB_OPERATION_NAME,
|
||||
DB_COLLECTION_NAME,
|
||||
PARAM_COUNT,
|
||||
PARAM_SETS,
|
||||
TIME_LIMIT_MS,
|
||||
ROWS_RETURNED,
|
||||
TRUNCATED,
|
||||
INTERRUPTED,
|
||||
SQL_ERROR_SUPPRESSED,
|
||||
EXECUTESCRIPT,
|
||||
EXECUTEMANY,
|
||||
),
|
||||
kind=SpanKind.CLIENT,
|
||||
)
|
||||
|
||||
DB_QUERY_EXECUTE = SpanName(
|
||||
"db.query.execute",
|
||||
"The read executing inside a SQL worker thread. Child of ``db.query``; the "
|
||||
"gap between the two is time spent waiting for a thread.",
|
||||
)
|
||||
|
||||
DB_WRITE_QUEUE_WAIT = SpanName(
|
||||
"db.write.queue_wait",
|
||||
"Time a write spent waiting in its database's write queue before the write "
|
||||
"thread picked it up. Child of ``db.query`` for a ``block=True`` write, "
|
||||
"where the caller awaits the write and containment is accurate. For a "
|
||||
"``block=False`` write the caller does not await it - the enqueueing "
|
||||
"request *caused* the write without *containing* it, and the write's "
|
||||
"spans can outlive the request's own - so this is a root span instead, "
|
||||
"carrying an OpenTelemetry link back to the enqueueing span rather than "
|
||||
"a parent. A link records causation without asserting containment, which "
|
||||
"is exactly the distinction here.",
|
||||
)
|
||||
|
||||
DB_WRITE_EXECUTE = SpanName(
|
||||
"db.write.execute",
|
||||
"The write executing on the write thread. Child of ``db.query`` for a "
|
||||
"``block=True`` write; for ``block=False`` a root span with a link back "
|
||||
"to the enqueueing span instead - see ``db.write.queue_wait`` above.",
|
||||
(ISOLATED_CONNECTION, TRANSACTION),
|
||||
)
|
||||
|
||||
STARTUP = SpanName(
|
||||
"datasette.startup",
|
||||
"``invoke_startup()`` running: ``register_events``, ``register_actions``, "
|
||||
"``register_column_types``, ``prepare_jinja2_environment``, internal-database "
|
||||
"schema catalog refresh (including the ``prepare_connection`` warm-up this "
|
||||
"triggers for each database touched for the first time), saved queries, "
|
||||
"column type config and the ``startup`` hook. Runs once per process, before "
|
||||
"any request exists, so without this span every child it creates would be "
|
||||
"its own orphan root trace. A connection warmed later - lazily, the first "
|
||||
"time a *request* touches a new database or thread - nests under that "
|
||||
"request's own span instead, not under this one, since this span has "
|
||||
"already ended by then.",
|
||||
)
|
||||
|
||||
SPANS = (
|
||||
HTTP_REQUEST,
|
||||
DB_QUERY,
|
||||
DB_QUERY_EXECUTE,
|
||||
DB_WRITE_QUEUE_WAIT,
|
||||
DB_WRITE_EXECUTE,
|
||||
STARTUP,
|
||||
)
|
||||
|
||||
|
||||
def span_for(emitted_name, kind=None, spans=None):
|
||||
"""
|
||||
Resolve an emitted span name to its registry entry, or None.
|
||||
|
||||
Handles the two entry kinds whose emitted names are not knowable in
|
||||
advance:
|
||||
|
||||
- `prefix=True` - the name carries a variable suffix after a fixed
|
||||
prefix, matched by prefix. Core registers none; plugin registries use
|
||||
it for names like ``chat {model}``.
|
||||
- `dynamic=True` - the name has no fixed part at all, so it is matched
|
||||
on `kind` instead and the caller has to supply one.
|
||||
|
||||
Exact matches win over prefix matches, and both win over dynamic, so a
|
||||
looser entry can never shadow a span with a registered name.
|
||||
|
||||
`spans` defaults to core's own registry; the plugin testing kit passes a
|
||||
plugin's tuple instead.
|
||||
"""
|
||||
if spans is None:
|
||||
spans = SPANS
|
||||
for span in spans:
|
||||
if span.dynamic:
|
||||
continue
|
||||
if emitted_name == span:
|
||||
return span
|
||||
for span in spans:
|
||||
if span.prefix and emitted_name.startswith(span):
|
||||
return span
|
||||
if kind is not None:
|
||||
for span in spans:
|
||||
if span.dynamic and span.kind == kind:
|
||||
return span
|
||||
return None
|
||||
|
||||
|
||||
def metric_for(emitted_name, metrics=None):
|
||||
"""
|
||||
Resolve an emitted metric name to its registry entry, or None.
|
||||
|
||||
The `span_for()` analogue - simpler, because metric names are always
|
||||
static strings. `metrics` defaults to core's own registry; the plugin
|
||||
testing kit passes a plugin's tuple instead.
|
||||
"""
|
||||
if metrics is None:
|
||||
metrics = METRICS
|
||||
for metric in metrics:
|
||||
if emitted_name == metric:
|
||||
return metric
|
||||
return None
|
||||
|
||||
|
||||
def attribute_allowed(entry, emitted_key):
|
||||
"""
|
||||
Whether `emitted_key` is a registered attribute of `entry`.
|
||||
|
||||
`entry` is a `SpanName` or a `MetricName` - both carry `.attributes`.
|
||||
"""
|
||||
if entry is None:
|
||||
return False
|
||||
return emitted_key in entry.attributes
|
||||
|
||||
|
||||
def attribute_value_allowed(entry, emitted_key, value):
|
||||
"""
|
||||
Whether `value` is permitted for `emitted_key` on `entry` (a `SpanName`
|
||||
or a `MetricName`).
|
||||
|
||||
True for any value when the attribute declares no `values=` enum; when it
|
||||
does, membership is enforced - that is what makes a declared enum a real
|
||||
cardinality bound rather than documentation. On a metric entry this is
|
||||
where the bound matters most: a metric series is keyed by its attribute
|
||||
values.
|
||||
"""
|
||||
if entry is None:
|
||||
return False
|
||||
for attribute in entry.attributes:
|
||||
if attribute == emitted_key:
|
||||
return attribute.values is None or value in attribute.values
|
||||
return False
|
||||
|
||||
|
||||
# --- Metrics --------------------------------------------------------------
|
||||
|
||||
# Every duration histogram here is in seconds, and OpenTelemetry's default
|
||||
# bucket boundaries are tuned for milliseconds - their first non-zero boundary
|
||||
# is 5, so without explicit boundaries every SQLite query lands in the single
|
||||
# (0, 5] second bucket and every quantile query returns noise.
|
||||
#
|
||||
# These are the OpenTelemetry semantic conventions' recommended boundaries for
|
||||
# db.client.operation.duration, in seconds, plus 0.0001 and 0.0005 at the
|
||||
# bottom. The deviation is deliberate: those boundaries assume a network
|
||||
# database client, whereas SQLite is in-process and a large fraction of real
|
||||
# queries run in 30-80us, which would otherwise all pile into the first
|
||||
# bucket and be indistinguishable from each other.
|
||||
#
|
||||
# One shared list is used for every duration histogram rather than a tailored
|
||||
# list each, so that dashboards stay comparable and a queue wait can be read
|
||||
# against the query duration it delays. It already spans 100us to 10s, which
|
||||
# covers both a fast in-process read and a write queued behind contention.
|
||||
DURATION_BUCKETS = (0.0001, 0.0005, 0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1, 5, 10)
|
||||
|
||||
M_OPERATION_DURATION = MetricName(
|
||||
"db.client.operation.duration",
|
||||
HISTOGRAM,
|
||||
"s",
|
||||
"Duration of a SQL operation. The standard OpenTelemetry semantic "
|
||||
"convention metric, and the one that survives trace sampling. "
|
||||
"Callback-style calls (``execute_fn()`` and friends) are counted "
|
||||
"alongside the SQL-string methods.",
|
||||
(DB_SYSTEM, DB_NAMESPACE, OPERATION, ERROR_TYPE),
|
||||
buckets=DURATION_BUCKETS,
|
||||
)
|
||||
|
||||
M_WRITE_QUEUE_WAIT = MetricName(
|
||||
"datasette.write.queue_wait",
|
||||
HISTOGRAM,
|
||||
"s",
|
||||
"Time each write waited in its database's write queue. The metric "
|
||||
"counterpart of the ``db.write.queue_wait`` span.",
|
||||
(DB_NAMESPACE,),
|
||||
buckets=DURATION_BUCKETS,
|
||||
)
|
||||
|
||||
M_QUERIES_INTERRUPTED = MetricName(
|
||||
"datasette.sql.queries.interrupted",
|
||||
COUNTER,
|
||||
"{query}",
|
||||
"Queries cancelled for exceeding :ref:`setting_sql_time_limit_ms`. Worth "
|
||||
"alerting on: a rising rate means the limit is too tight or a table has "
|
||||
"outgrown its queries. A caller that opted into a deliberately shorter "
|
||||
"budget - facet suggestion, for example - is not counted, for the same "
|
||||
"reason its timeout is not a span error.",
|
||||
(DB_NAMESPACE,),
|
||||
)
|
||||
|
||||
M_THREADS_LIMIT = MetricName(
|
||||
"datasette.sql.threads.limit",
|
||||
GAUGE,
|
||||
"{thread}",
|
||||
"Maximum concurrent read queries - the :ref:`setting_num_sql_threads` "
|
||||
"value. Not reported when ``num_sql_threads`` is ``0``, since then queries "
|
||||
"run on the event loop and there is no pool.",
|
||||
)
|
||||
|
||||
M_THREADS_QUEUE_DEPTH = MetricName(
|
||||
"datasette.sql.threads.queue_depth",
|
||||
GAUGE,
|
||||
"{query}",
|
||||
"Read queries waiting for a free thread. **This is the saturation "
|
||||
"signal** - sustained above zero means requests are queueing on "
|
||||
"``num_sql_threads``.",
|
||||
)
|
||||
|
||||
M_QUERIES_PENDING = MetricName(
|
||||
"datasette.sql.queries.pending",
|
||||
GAUGE,
|
||||
"{query}",
|
||||
"Read queries submitted to the pool and not yet complete. Summed across "
|
||||
"databases and compared against the thread limit, this is pool "
|
||||
"utilisation.",
|
||||
(DB_NAMESPACE,),
|
||||
)
|
||||
|
||||
M_WRITE_QUEUE_DEPTH = MetricName(
|
||||
"datasette.write.queue_depth",
|
||||
GAUGE,
|
||||
"{write}",
|
||||
"Writes queued behind a database's single write thread. Backpressure that "
|
||||
"raising ``num_sql_threads`` cannot relieve. Not reported for a database "
|
||||
"that has never been written to.",
|
||||
(DB_NAMESPACE,),
|
||||
)
|
||||
|
||||
M_CONNECTIONS_OPEN = MetricName(
|
||||
"datasette.connections.open",
|
||||
GAUGE,
|
||||
"{connection}",
|
||||
"Open SQLite file connections currently tracked for closing.",
|
||||
(DB_NAMESPACE,),
|
||||
)
|
||||
|
||||
METRICS = (
|
||||
M_OPERATION_DURATION,
|
||||
M_WRITE_QUEUE_WAIT,
|
||||
M_QUERIES_INTERRUPTED,
|
||||
M_THREADS_LIMIT,
|
||||
M_THREADS_QUEUE_DEPTH,
|
||||
M_QUERIES_PENDING,
|
||||
M_WRITE_QUEUE_DEPTH,
|
||||
M_CONNECTIONS_OPEN,
|
||||
)
|
||||
|
|
@ -1,509 +0,0 @@
|
|||
"""
|
||||
Pytest helpers for testing OpenTelemetry instrumentation - Datasette's own
|
||||
and any plugin's. Part of Datasette's public plugin API; see the "Telemetry
|
||||
for plugin authors" documentation.
|
||||
|
||||
Usage from a plugin's ``conftest.py``::
|
||||
|
||||
from datasette.telemetry_testing import ( # noqa: F401
|
||||
MetricsCollector,
|
||||
otel_metrics,
|
||||
otel_meter_provider,
|
||||
otel_provider,
|
||||
otel_spans,
|
||||
)
|
||||
|
||||
Importing the fixture names into a conftest registers them; ``otel_provider``
|
||||
and ``otel_meter_provider`` are session-scoped and autouse, so a real SDK
|
||||
provider (when the SDK is installed) is in place before any test emits a
|
||||
signal. Tests then take ``otel_spans`` / ``otel_metrics``. Everything here
|
||||
imports the OpenTelemetry SDK lazily: with no SDK installed the fixtures
|
||||
skip rather than fail, and importing this module costs nothing.
|
||||
|
||||
The conformance helpers (`assert_spans_conform`, `assert_spans_covered`)
|
||||
check a registry of `SpanName` entries against actually-finished spans in
|
||||
both directions - emitted-but-unregistered and registered-but-never-emitted,
|
||||
the two drift modes documented in `tests/test_telemetry_registry.py`.
|
||||
"""
|
||||
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
import pytest
|
||||
|
||||
from .telemetry_registry import (
|
||||
attribute_allowed,
|
||||
attribute_value_allowed,
|
||||
metric_for,
|
||||
span_for,
|
||||
)
|
||||
|
||||
_span_exporter = None
|
||||
_metric_reader = None
|
||||
|
||||
|
||||
def install_span_exporter():
|
||||
"""
|
||||
Install a TracerProvider + InMemorySpanExporter once per process and
|
||||
return the exporter, or None when the SDK is not installed.
|
||||
|
||||
`set_tracer_provider()` is effectively once-per-process (a second call
|
||||
logs a warning and is ignored), so this must run before anything asserts
|
||||
on spans. A `SimpleSpanProcessor` exports synchronously on span end - no
|
||||
background batching thread, so assertions immediately after a request
|
||||
never race.
|
||||
"""
|
||||
global _span_exporter
|
||||
if _span_exporter is not None:
|
||||
return _span_exporter
|
||||
try:
|
||||
from opentelemetry import trace as otel_trace
|
||||
from opentelemetry.sdk.trace import TracerProvider
|
||||
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
|
||||
from opentelemetry.sdk.trace.export.in_memory_span_exporter import (
|
||||
InMemorySpanExporter,
|
||||
)
|
||||
except ImportError:
|
||||
return None
|
||||
exporter = InMemorySpanExporter()
|
||||
provider = TracerProvider()
|
||||
provider.add_span_processor(SimpleSpanProcessor(exporter))
|
||||
otel_trace.set_tracer_provider(provider)
|
||||
# set_tracer_provider() is once-per-process: if something else installed
|
||||
# a provider first (another conftest, opentelemetry-instrument, an
|
||||
# embedding app), the call above was silently ignored - and an exporter
|
||||
# wired to nothing would make every span assertion fail confusingly, or
|
||||
# pass vacuously on empty input. Leave the global unset in that case so
|
||||
# the fixtures skip with a clear message instead.
|
||||
if otel_trace.get_tracer_provider() is not provider:
|
||||
return None
|
||||
_span_exporter = exporter
|
||||
return exporter
|
||||
|
||||
|
||||
def install_metric_reader():
|
||||
"""
|
||||
Install a MeterProvider + InMemoryMetricReader once per process and
|
||||
return the reader, or None when the SDK is not installed.
|
||||
|
||||
DELTA temporality for counters and histograms, so each collection
|
||||
reports only what happened since the previous one - with the SDK default
|
||||
of CUMULATIVE, every metrics test would see every measurement from every
|
||||
earlier test in the session.
|
||||
"""
|
||||
global _metric_reader
|
||||
if _metric_reader is not None:
|
||||
return _metric_reader
|
||||
try:
|
||||
from opentelemetry import metrics as otel_metrics_api
|
||||
from opentelemetry.sdk.metrics import Counter, Histogram, MeterProvider
|
||||
from opentelemetry.sdk.metrics.export import (
|
||||
AggregationTemporality,
|
||||
InMemoryMetricReader,
|
||||
)
|
||||
except ImportError:
|
||||
return None
|
||||
reader = InMemoryMetricReader(
|
||||
preferred_temporality={
|
||||
Counter: AggregationTemporality.DELTA,
|
||||
Histogram: AggregationTemporality.DELTA,
|
||||
}
|
||||
)
|
||||
provider = MeterProvider(metric_readers=[reader])
|
||||
otel_metrics_api.set_meter_provider(provider)
|
||||
# Same once-per-process guard as the tracer side: a provider that did
|
||||
# not take must not leave a reader that collects nothing.
|
||||
if otel_metrics_api.get_meter_provider() is not provider:
|
||||
return None
|
||||
_metric_reader = reader
|
||||
return reader
|
||||
|
||||
|
||||
@pytest.fixture(scope="session", autouse=True)
|
||||
def otel_provider():
|
||||
"""
|
||||
Session-scoped, autouse: install the span exporter exactly once, before
|
||||
any span is created.
|
||||
|
||||
`datasette.telemetry.tracer` (and a plugin's own tracer) is a
|
||||
module-level `ProxyTracer`: once a provider exists, the first span it
|
||||
starts resolves a concrete tracer and caches it permanently. It does
|
||||
*not* cache the no-op tracer, so a span started before this fixture runs
|
||||
is merely lost rather than poisoning the tracer for the process. With no
|
||||
SDK installed this does nothing and spans stay no-op.
|
||||
"""
|
||||
install_span_exporter()
|
||||
|
||||
|
||||
@pytest.fixture(scope="session", autouse=True)
|
||||
def otel_meter_provider():
|
||||
"""
|
||||
Session-scoped, autouse: install the metric reader once per process.
|
||||
|
||||
Unlike the tracer, ordering is not load-bearing - `_ProxyMeter` and its
|
||||
instruments forward to a provider installed after they were created.
|
||||
Still autouse for symmetry, and so a single reader collects all run.
|
||||
"""
|
||||
install_metric_reader()
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def otel_reset():
|
||||
"""
|
||||
Autouse, function-scoped: drain the span exporter and metric reader
|
||||
after every test - including the ones that never look at telemetry.
|
||||
|
||||
Without this, every test that exercises the app leaves its recorded
|
||||
spans in the session-scoped exporter's list forever: a large suite
|
||||
accumulates hundreds of thousands of ReadableSpans, degrading memory
|
||||
and per-span export cost as the run goes on. Draining the metric reader
|
||||
likewise stops delta state piling up between metric tests.
|
||||
"""
|
||||
yield
|
||||
if _span_exporter is not None:
|
||||
_span_exporter.clear()
|
||||
if _metric_reader is not None:
|
||||
_metric_reader.get_metrics_data()
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def otel_spans():
|
||||
"""
|
||||
Function-scoped access to the finished-spans exporter: clears spans left
|
||||
over from previous tests, then yields the exporter so a test can call
|
||||
`.get_finished_spans()`. Skips if the OTel SDK is not installed.
|
||||
"""
|
||||
pytest.importorskip("opentelemetry.sdk")
|
||||
exporter = install_span_exporter()
|
||||
if exporter is None:
|
||||
pytest.skip("OpenTelemetry SDK provider was not installed")
|
||||
exporter.clear()
|
||||
yield exporter
|
||||
|
||||
|
||||
class MetricsCollector:
|
||||
"""
|
||||
Thin reader over an `InMemoryMetricReader`.
|
||||
|
||||
`collect()` runs a collection cycle - which is what invokes observable
|
||||
gauge callbacks - and snapshots the result. Queries then run against
|
||||
that snapshot rather than re-collecting, so a test that inspects
|
||||
several metrics sees one consistent moment and does not drain delta
|
||||
state twice.
|
||||
"""
|
||||
|
||||
def __init__(self, reader):
|
||||
self.reader = reader
|
||||
self.snapshot = {}
|
||||
# (instrumentation scope name, sdk Metric) pairs from the last
|
||||
# collect() - the metric conformance helpers read this, because the
|
||||
# name-keyed snapshot deliberately flattens the scope away.
|
||||
self.collected = []
|
||||
|
||||
def collect(self):
|
||||
self.snapshot = {}
|
||||
self.collected = []
|
||||
data = self.reader.get_metrics_data()
|
||||
if data is None:
|
||||
return self.snapshot
|
||||
for resource_metrics in data.resource_metrics:
|
||||
for scope_metrics in resource_metrics.scope_metrics:
|
||||
scope_name = scope_metrics.scope.name if scope_metrics.scope else None
|
||||
for metric in scope_metrics.metrics:
|
||||
self.snapshot.setdefault(metric.name, []).extend(
|
||||
metric.data.data_points
|
||||
)
|
||||
self.collected.append((scope_name, metric))
|
||||
return self.snapshot
|
||||
|
||||
def points(self, name, attributes=None):
|
||||
"Data points for `name` whose attributes are a superset of `attributes`."
|
||||
found = []
|
||||
for point in self.snapshot.get(name, []):
|
||||
point_attributes = dict(point.attributes or {})
|
||||
if all(point_attributes.get(k) == v for k, v in (attributes or {}).items()):
|
||||
found.append(point)
|
||||
return found
|
||||
|
||||
def point(self, name, attributes=None):
|
||||
"The single matching data point, asserting there is exactly one."
|
||||
found = self.points(name, attributes)
|
||||
assert len(found) == 1, (
|
||||
f"expected exactly one {name} point matching {attributes}, "
|
||||
f"got {len(found)}: {found}"
|
||||
)
|
||||
return found[0]
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def otel_metrics():
|
||||
"""
|
||||
Function-scoped metrics collector. Drains delta state accumulated by
|
||||
earlier tests before yielding, so counts start from zero.
|
||||
"""
|
||||
pytest.importorskip("opentelemetry.sdk")
|
||||
reader = install_metric_reader()
|
||||
if reader is None:
|
||||
pytest.skip("OpenTelemetry SDK meter provider was not installed")
|
||||
reader.get_metrics_data()
|
||||
yield MetricsCollector(reader)
|
||||
|
||||
|
||||
def _scoped(finished_spans, scope_name):
|
||||
if scope_name is None:
|
||||
return list(finished_spans)
|
||||
return [
|
||||
span
|
||||
for span in finished_spans
|
||||
if span.instrumentation_scope and span.instrumentation_scope.name == scope_name
|
||||
]
|
||||
|
||||
|
||||
def assert_spans_conform(registry_spans, finished_spans, scope_name=None):
|
||||
"""
|
||||
Every finished span (optionally: only those from `scope_name`, which is
|
||||
what a plugin should pass - its own tracer's name) resolves to an entry
|
||||
in `registry_spans`, sets only registered attributes, and respects any
|
||||
declared `values=` enums. This is the emitted-but-unregistered direction:
|
||||
instrumentation added without documentation fails here.
|
||||
"""
|
||||
problems = []
|
||||
for span in _scoped(finished_spans, scope_name):
|
||||
entry = span_for(str(span.name), kind=span.kind, spans=registry_spans)
|
||||
if entry is None:
|
||||
problems.append(f"unregistered span: {span.name!r}")
|
||||
continue
|
||||
for key, value in (span.attributes or {}).items():
|
||||
if not attribute_allowed(entry, str(key)):
|
||||
problems.append(f"{span.name}: unregistered attribute {key!r}")
|
||||
elif not attribute_value_allowed(entry, str(key), value):
|
||||
problems.append(
|
||||
f"{span.name}: {key}={value!r} not in the declared enum"
|
||||
)
|
||||
assert not problems, "\n".join(problems)
|
||||
|
||||
|
||||
def assert_spans_covered(registry_spans, finished_spans, scope_name=None):
|
||||
"""
|
||||
Every entry in `registry_spans` was emitted at least once, and every one
|
||||
of its registered non-`optional` attributes appeared on it at least
|
||||
once. This is the registered-but-never-emitted direction - documentation
|
||||
describing a signal that no longer exists, which is worse than omitting
|
||||
it because a reader will build a dashboard on it. Run it against a
|
||||
workload broad enough to exercise everything the registry claims;
|
||||
`optional=True` attributes are exempt so a workload is not forced to
|
||||
manufacture every error path (pin those with targeted tests instead).
|
||||
"""
|
||||
spans = _scoped(finished_spans, scope_name)
|
||||
seen_attributes = {}
|
||||
for span in spans:
|
||||
entry = span_for(str(span.name), kind=span.kind, spans=registry_spans)
|
||||
if entry is not None:
|
||||
seen = seen_attributes.setdefault(str(entry), set())
|
||||
seen.update(str(key) for key in (span.attributes or {}))
|
||||
problems = []
|
||||
for entry in registry_spans:
|
||||
if str(entry) not in seen_attributes:
|
||||
problems.append(f"registered span never emitted: {entry!r}")
|
||||
continue
|
||||
required = {
|
||||
str(attribute) for attribute in entry.attributes if not attribute.optional
|
||||
}
|
||||
missing = required - seen_attributes[str(entry)]
|
||||
if missing:
|
||||
problems.append(
|
||||
f"{entry}: registered attributes never emitted: {sorted(missing)}"
|
||||
)
|
||||
assert not problems, "\n".join(problems)
|
||||
|
||||
|
||||
# Registry instrument kinds mapped to the SDK data type collected for them.
|
||||
# A registry kind outside this table (a plugin's own vocabulary) is not
|
||||
# kind-checked. Both counter kinds collect as Sum; monotonicity is what
|
||||
# tells them apart, checked separately below.
|
||||
_KIND_TO_DATA_TYPE = {
|
||||
"Counter": "Sum",
|
||||
"UpDownCounter": "Sum",
|
||||
"Histogram": "Histogram",
|
||||
"Observable gauge": "Gauge",
|
||||
}
|
||||
_KIND_IS_MONOTONIC = {"Counter": True, "UpDownCounter": False}
|
||||
|
||||
|
||||
def _scoped_metrics(collector, scope_name):
|
||||
for scope, metric in collector.collected:
|
||||
if scope_name is None or scope == scope_name:
|
||||
yield metric
|
||||
|
||||
|
||||
def assert_metrics_conform(registry_metrics, collector, scope_name=None):
|
||||
"""
|
||||
Every metric in the collector's last `collect()` (optionally: only those
|
||||
from `scope_name`, which is what a plugin should pass - its own meter's
|
||||
name) is registered in `registry_metrics`, was created as the instrument
|
||||
kind and unit the registry declares, sets only registered attributes,
|
||||
and respects any declared `values=` enums.
|
||||
|
||||
The kind and unit checks catch a drift nothing else does: the registry
|
||||
entry and the `meter.create_*()` call are separate statements, and a
|
||||
dashboard built on the registry's word breaks silently if they disagree.
|
||||
"""
|
||||
problems = set()
|
||||
for metric in _scoped_metrics(collector, scope_name):
|
||||
entry = metric_for(metric.name, metrics=registry_metrics)
|
||||
if entry is None:
|
||||
problems.add(f"unregistered metric: {metric.name!r}")
|
||||
continue
|
||||
expected_data_type = _KIND_TO_DATA_TYPE.get(entry.kind)
|
||||
actual_data_type = type(metric.data).__name__
|
||||
if expected_data_type is not None and actual_data_type != expected_data_type:
|
||||
problems.add(
|
||||
f"{metric.name}: registry declares {entry.kind}, "
|
||||
f"SDK collected {actual_data_type}"
|
||||
)
|
||||
expected_monotonic = _KIND_IS_MONOTONIC.get(entry.kind)
|
||||
actual_monotonic = getattr(metric.data, "is_monotonic", None)
|
||||
if (
|
||||
expected_monotonic is not None
|
||||
and actual_monotonic is not None
|
||||
and actual_monotonic != expected_monotonic
|
||||
):
|
||||
problems.add(
|
||||
f"{metric.name}: registry declares {entry.kind}, but the "
|
||||
f"collected Sum is_monotonic={actual_monotonic}"
|
||||
)
|
||||
if (metric.unit or "") != (entry.unit or ""):
|
||||
problems.add(
|
||||
f"{metric.name}: instrument unit {metric.unit!r} != "
|
||||
f"registry unit {entry.unit!r}"
|
||||
)
|
||||
for point in metric.data.data_points:
|
||||
for key, value in dict(point.attributes or {}).items():
|
||||
if not attribute_allowed(entry, str(key)):
|
||||
problems.add(f"{metric.name}: unregistered attribute {key!r}")
|
||||
elif not attribute_value_allowed(entry, str(key), value):
|
||||
problems.add(
|
||||
f"{metric.name}: {key}={value!r} not in the declared enum"
|
||||
)
|
||||
assert not problems, "\n".join(sorted(problems))
|
||||
|
||||
|
||||
def assert_metrics_covered(registry_metrics, collector, scope_name=None):
|
||||
"""
|
||||
Every entry in `registry_metrics` was collected at least once, and every
|
||||
registered non-`optional` attribute appeared on it at least once - the
|
||||
registered-but-never-emitted direction for metrics.
|
||||
|
||||
Run one broad workload, then a single `collect()`, then this: the reader
|
||||
uses delta temporality, so measurements drained by an earlier collect()
|
||||
are gone. `optional=True` attributes (e.g. an `error.type` only present
|
||||
on failures) are exempt, same as the span-side helper.
|
||||
"""
|
||||
seen_attributes = {}
|
||||
for metric in _scoped_metrics(collector, scope_name):
|
||||
entry = metric_for(metric.name, metrics=registry_metrics)
|
||||
if entry is None:
|
||||
continue
|
||||
seen = seen_attributes.setdefault(str(entry), set())
|
||||
for point in metric.data.data_points:
|
||||
seen.update(str(key) for key in dict(point.attributes or {}))
|
||||
problems = []
|
||||
for entry in registry_metrics:
|
||||
if str(entry) not in seen_attributes:
|
||||
problems.append(f"registered metric never collected: {entry!r}")
|
||||
continue
|
||||
required = {
|
||||
str(attribute) for attribute in entry.attributes if not attribute.optional
|
||||
}
|
||||
missing = required - seen_attributes[str(entry)]
|
||||
if missing:
|
||||
problems.append(
|
||||
f"{entry}: registered attributes never collected: {sorted(missing)}"
|
||||
)
|
||||
assert not problems, "\n".join(problems)
|
||||
|
||||
|
||||
def assert_no_forbidden_values(
|
||||
forbidden, finished_spans=None, collector=None, scope_name=None
|
||||
):
|
||||
"""
|
||||
Assert that none of the `forbidden` strings appear anywhere in the
|
||||
emitted telemetry: span names, span attribute values, span event names
|
||||
and attributes, span status descriptions, or metric point attributes.
|
||||
|
||||
This is the enforcement half of the privacy rules in the plugin
|
||||
telemetry documentation. The strongest way to use it is to *plant*
|
||||
sentinel values in your test workload - a fake email address, a token,
|
||||
a username your fixtures log in with - and assert they never leak into
|
||||
a signal:
|
||||
|
||||
FORBIDDEN = {"secret-token-123", "alice@example.com"}
|
||||
run_workload_using_those_values()
|
||||
assert_no_forbidden_values(
|
||||
FORBIDDEN,
|
||||
finished_spans=otel_spans.get_finished_spans(),
|
||||
collector=otel_metrics,
|
||||
scope_name="my_plugin",
|
||||
)
|
||||
|
||||
Matching is plain substring on the string form of each value; empty
|
||||
strings in `forbidden` are ignored. Pass `finished_spans` and/or a
|
||||
collected `MetricsCollector`; `scope_name=None` checks every scope,
|
||||
which is the right default here - a leak through *core's* signals (e.g.
|
||||
SQL text carrying a secret) is still a leak.
|
||||
"""
|
||||
needles = [needle for needle in forbidden if needle]
|
||||
leaks = set()
|
||||
|
||||
def check(value, where):
|
||||
text = str(value)
|
||||
for needle in needles:
|
||||
if needle in text:
|
||||
leaks.add(f"{where} contains {needle!r}")
|
||||
|
||||
if finished_spans is not None:
|
||||
for span in _scoped(finished_spans, scope_name):
|
||||
check(span.name, f"span name {str(span.name)!r}")
|
||||
for key, value in (span.attributes or {}).items():
|
||||
check(value, f"{span.name} attribute {key}")
|
||||
for event in span.events or ():
|
||||
check(event.name, f"{span.name} event name")
|
||||
for key, value in (event.attributes or {}).items():
|
||||
check(value, f"{span.name} event {event.name} attribute {key}")
|
||||
if span.status is not None and span.status.description:
|
||||
check(span.status.description, f"{span.name} status description")
|
||||
if collector is not None:
|
||||
for metric in _scoped_metrics(collector, scope_name):
|
||||
for point in metric.data.data_points:
|
||||
for key, value in dict(point.attributes or {}).items():
|
||||
check(value, f"metric {metric.name} attribute {key}")
|
||||
assert not leaks, "forbidden values leaked into telemetry:\n" + "\n".join(
|
||||
sorted(leaks)
|
||||
)
|
||||
|
||||
|
||||
def assert_package_never_imports_sdk(*module_names):
|
||||
"""
|
||||
Import the named modules in a fresh interpreter and assert none of them
|
||||
dragged in `opentelemetry.sdk`. Checked via sys.modules in a subprocess
|
||||
rather than by grepping, so a lazy `import opentelemetry.sdk` inside a
|
||||
function body cannot slip past. A plugin should depend on
|
||||
`opentelemetry-api` only, exactly as Datasette core does.
|
||||
|
||||
Run the test that calls this early in your suite: on macOS/CPython 3.13
|
||||
a process that has accumulated many threads can crash (SIGBUS) in
|
||||
subprocess's fork+exec - Datasette's own conftest front-loads its
|
||||
equivalent tests by name for exactly this reason.
|
||||
"""
|
||||
imports = "; ".join(f"import {name}" for name in module_names)
|
||||
code = (
|
||||
f"import sys; {imports}; "
|
||||
"print([m for m in sys.modules if m.startswith('opentelemetry.sdk')])"
|
||||
)
|
||||
result = subprocess.run(
|
||||
[sys.executable, "-c", code], capture_output=True, text=True, check=True
|
||||
)
|
||||
assert result.stdout.strip() == "[]", (
|
||||
f"importing {module_names} pulled in the OpenTelemetry SDK: "
|
||||
f"{result.stdout.strip()}"
|
||||
)
|
||||
|
|
@ -407,7 +407,7 @@ class RowView(BaseView):
|
|||
raise Forbidden("You do not have permission to view this table")
|
||||
|
||||
results = await resolved.db.execute(
|
||||
resolved.sql, resolved.params, truncate=True, table=table
|
||||
resolved.sql, resolved.params, truncate=True
|
||||
)
|
||||
columns = [r[0] for r in results.description]
|
||||
rows = list(results.rows)
|
||||
|
|
@ -652,9 +652,6 @@ class RowView(BaseView):
|
|||
]
|
||||
)
|
||||
try:
|
||||
# No table= here: this counts incoming references across every
|
||||
# foreign key pointing at this row, so it spans many tables and
|
||||
# there is no single value db.collection.name could take.
|
||||
rows = list(await db.execute(sql, {"id": pk_values[0]}))
|
||||
except QueryInterrupted:
|
||||
# Almost certainly hit the timeout
|
||||
|
|
@ -843,7 +840,7 @@ class RowUpdateView(BaseView):
|
|||
returned_row = None
|
||||
if data.get("return"):
|
||||
results = await resolved.db.execute(
|
||||
resolved.sql, resolved.params, truncate=True, table=resolved.table
|
||||
resolved.sql, resolved.params, truncate=True
|
||||
)
|
||||
returned_row = results.dicts()[0]
|
||||
result["rows"] = [returned_row]
|
||||
|
|
@ -861,7 +858,7 @@ class RowUpdateView(BaseView):
|
|||
message_row = returned_row
|
||||
if message_row is None:
|
||||
results = await resolved.db.execute(
|
||||
resolved.sql, resolved.params, truncate=True, table=resolved.table
|
||||
resolved.sql, resolved.params, truncate=True
|
||||
)
|
||||
message_row = results.first()
|
||||
self.ds.add_message(
|
||||
|
|
|
|||
|
|
@ -1170,7 +1170,6 @@ class TableInsertView(BaseView):
|
|||
"rowid, " if pks == ["rowid"] else "", table_name, where_clause
|
||||
),
|
||||
args,
|
||||
table=table_name,
|
||||
)
|
||||
result["rows"] = fetched_rows.dicts()
|
||||
else:
|
||||
|
|
@ -1383,9 +1382,7 @@ class TableDropView(BaseView):
|
|||
"database": database_name,
|
||||
"table": table_name,
|
||||
"row_count": (
|
||||
await db.execute(
|
||||
f"select count(*) from [{table_name}]", table=table_name
|
||||
)
|
||||
await db.execute(f"select count(*) from [{table_name}]")
|
||||
).single_value(),
|
||||
"message": 'Pass "confirm": true to confirm',
|
||||
},
|
||||
|
|
@ -1579,10 +1576,7 @@ class TableAutocompleteView(BaseView):
|
|||
|
||||
try:
|
||||
results = await db.execute(
|
||||
sql,
|
||||
params,
|
||||
custom_time_limit=AUTOCOMPLETE_TIME_LIMIT_MS,
|
||||
table=table_name,
|
||||
sql, params, custom_time_limit=AUTOCOMPLETE_TIME_LIMIT_MS
|
||||
)
|
||||
except QueryInterrupted:
|
||||
fallback_where = _autocomplete_prefix_like(pks[0])
|
||||
|
|
@ -1603,7 +1597,6 @@ class TableAutocompleteView(BaseView):
|
|||
fallback_sql,
|
||||
params,
|
||||
custom_time_limit=AUTOCOMPLETE_TIME_LIMIT_MS,
|
||||
table=table_name,
|
||||
)
|
||||
except QueryInterrupted:
|
||||
return Response.json({"ok": True, "rows": []})
|
||||
|
|
@ -2170,9 +2163,7 @@ async def table_view_data(
|
|||
|
||||
# Execute the main query!
|
||||
try:
|
||||
results = await db.execute(
|
||||
sql, params, truncate=True, table=table_name, **extra_args
|
||||
)
|
||||
results = await db.execute(sql, params, truncate=True, **extra_args)
|
||||
except (sqlite3.OperationalError, InvalidSql) as e:
|
||||
raise DatasetteError(str(e), title="Invalid SQL", status=400)
|
||||
|
||||
|
|
@ -2448,7 +2439,6 @@ async def _next_value_and_url(
|
|||
await db.execute(
|
||||
prefix_lookup_sql,
|
||||
{**{f"pk{i}": rows[-2][pk] for i, pk in enumerate(pks)}},
|
||||
table=table_name,
|
||||
)
|
||||
).single_value()
|
||||
if isinstance(prefix, dict) and "value" in prefix:
|
||||
|
|
|
|||
|
|
@ -4,20 +4,6 @@
|
|||
Changelog
|
||||
=========
|
||||
|
||||
.. _v_unreleased:
|
||||
|
||||
Unreleased
|
||||
----------
|
||||
|
||||
- Datasette's database layer now emits `OpenTelemetry <https://opentelemetry.io/>`__ spans: one per query, covering the full round trip including time spent waiting for a SQL worker thread, plus separate child spans for the execution itself and for time spent in the write queue. Callback-style calls - :ref:`db.execute_fn() <database_execute_fn>`, :ref:`db.execute_write_fn() <database_execute_write_fn>` and ``db.execute_isolated_fn()``, the documented way for plugins to run arbitrary SQL - are covered too, carrying ``datasette.callback`` in place of the SQL text. Datasette core depends on ``opentelemetry-api`` only and never installs an SDK provider, an exporter or a sampler, so there is no effect and no measurable overhead unless tracing is switched on externally - normally with the standard ``opentelemetry-instrument`` agent. See :ref:`internals_telemetry`. (:issue:`1730`)
|
||||
- :ref:`db.execute(sql, ..., table=None) <database_execute>` has a new optional ``table=`` parameter, naming the table a query is about so it can be recorded on that query's OpenTelemetry span. It has no effect on query execution, and Datasette never derives it from the SQL. (:issue:`1730`)
|
||||
- Every HTTP request now gets an OpenTelemetry ``SERVER`` span, named after the request method and matched route, carrying ``http.route``, the response status and W3C trace context extracted from inbound headers - so every database span has a request to belong to, and Datasette joins distributed traces started by a proxy or calling service. The query string is never recorded. See :ref:`internals_telemetry_requests`. (:issue:`1730`)
|
||||
- Datasette core now also emits OpenTelemetry **metrics** covering SQL thread pool saturation, per-database write queue depth, open connections, query latency and time-limit interruptions. These answer operational questions that spans structurally cannot - "am I saturating my :ref:`setting_num_sql_threads` threads?" is a level, not an event - and they survive trace sampling. As with spans, core installs no ``MeterProvider``, so there is no cost unless metrics are collected externally. See :ref:`internals_telemetry`. (:issue:`1730`)
|
||||
|
||||
- New :ref:`plugin telemetry kit <plugin_telemetry>` for plugins that emit their own OpenTelemetry signals: the registry classes (``Attribute`` with closed-enum ``values=``, ``SpanName`` with prefix-matched families, ``MetricName``) are now documented public API, ``datasette.telemetry.linked_root_span_kwargs()`` provides the root-span-with-link shape for background work, ``datasette.telemetry.request_span()`` is documented, and ``datasette.telemetry_testing`` ships the pytest fixtures and two-way conformance checks - for spans and metrics, including instrument kind/unit verification, enum enforcement and a forbidden-values privacy walk - that core's own suite uses. (:issue:`1730`)
|
||||
|
||||
Nothing is removed by the OpenTelemetry work: the ``?_trace=1`` query string parameter, the ``trace_debug`` setting and the :ref:`internals_tracer` module all continue to work as before.
|
||||
|
||||
.. _v1_0_a38:
|
||||
|
||||
1.0a38 (2026-08-06)
|
||||
|
|
|
|||
|
|
@ -64,7 +64,6 @@ Contents
|
|||
javascript_plugins
|
||||
plugin_hooks
|
||||
testing_plugins
|
||||
plugin_telemetry
|
||||
internals
|
||||
events
|
||||
upgrade_guide
|
||||
|
|
|
|||
|
|
@ -1963,9 +1963,6 @@ Executes a SQL query against the database and returns the resulting rows (see :r
|
|||
``log_sql_errors`` - boolean
|
||||
Should any SQL errors be logged to the console in addition to being raised as an error? Defaults to ``True``.
|
||||
|
||||
``table`` - string
|
||||
The name of the table this query is about, if the caller already knows it. This has no effect on how the query executes - it is recorded as the ``db.collection.name`` attribute on the :ref:`OpenTelemetry span <internals_telemetry>` for the query. Datasette never derives this from the SQL, so leave it unset for queries that do not have one obvious table.
|
||||
|
||||
.. _database_results:
|
||||
|
||||
Results
|
||||
|
|
@ -2024,8 +2021,6 @@ Example usage:
|
|||
|
||||
version = await db.execute_fn(get_version)
|
||||
|
||||
The call is traced as a ``db.query`` OpenTelemetry span carrying ``datasette.callback`` (the function's qualified name) rather than ``db.query.text``, since the SQL is whatever the function chooses to run - see :ref:`internals_telemetry`. Passing a named function gives the span a readable identity; a lambda reports ``<lambda>``.
|
||||
|
||||
.. _database_execute_write:
|
||||
|
||||
await db.execute_write(sql, params=None, block=True, request=None, return_all=False, returning_limit=10, transaction=True)
|
||||
|
|
@ -2102,8 +2097,6 @@ This method works like ``.execute_write()``, but instead of a SQL statement you
|
|||
|
||||
The function can then perform multiple actions, safe in the knowledge that it has exclusive access to the single writable connection for as long as it is executing.
|
||||
|
||||
Like ``execute_fn()``, the call is traced as a ``db.query`` OpenTelemetry span carrying ``datasette.callback`` rather than ``db.query.text``, above the write-queue spans - see :ref:`internals_telemetry`. A named function gives the span a readable identity; a lambda reports ``<lambda>``.
|
||||
|
||||
.. warning::
|
||||
|
||||
``fn`` needs to be a regular function, not an ``async def`` function.
|
||||
|
|
@ -2320,293 +2313,6 @@ The ``Database`` class also provides properties and methods for introspecting th
|
|||
}
|
||||
}
|
||||
|
||||
.. _internals_telemetry:
|
||||
|
||||
OpenTelemetry
|
||||
=============
|
||||
|
||||
Datasette core depends on `opentelemetry-api <https://pypi.org/project/opentelemetry-api/>`__ only. It never creates a ``TracerProvider``, never configures an exporter and never sets a sampler. With no OpenTelemetry SDK provider installed, every span described below is a no-op ``NonRecordingSpan``: nothing is recorded, nothing is exported, and the cost does not show up in page latency. Benchmarking a table page with and without this instrumentation, the median moved by less than the run-to-run variation of the benchmark itself.
|
||||
|
||||
Turning tracing on is entirely an operational decision made outside of Datasette itself: run Datasette under the standard ``opentelemetry-instrument`` agent, or embed Datasette inside a host application that installs its own provider.
|
||||
|
||||
Plugins can emit their own spans and metrics alongside these, using the same registry classes and test helpers core uses - see :ref:`plugin_telemetry`.
|
||||
|
||||
Everything Datasette emits carries the instrumentation scope ``datasette``, versioned with the running Datasette version and declaring the `semantic conventions schema <https://opentelemetry.io/docs/specs/otel/schemas/>`__ its attribute names follow.
|
||||
|
||||
This is separate from, and does not replace, the built-in :ref:`internals_tracer` mechanism behind ``?_trace=1`` and the :ref:`setting_trace_debug` setting. Both continue to work exactly as before.
|
||||
|
||||
.. _internals_telemetry_turning_on:
|
||||
|
||||
Turning tracing on
|
||||
------------------
|
||||
|
||||
Install an OpenTelemetry SDK, an exporter and the instrumentation agent, then launch Datasette through ``opentelemetry-instrument``:
|
||||
|
||||
.. code-block:: bash
|
||||
|
||||
pip install opentelemetry-distro opentelemetry-exporter-otlp
|
||||
OTEL_SERVICE_NAME=datasette \
|
||||
OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 \
|
||||
OTEL_METRICS_EXPORTER=none \
|
||||
OTEL_LOGS_EXPORTER=none \
|
||||
opentelemetry-instrument datasette mydb.db
|
||||
|
||||
Point ``OTEL_EXPORTER_OTLP_ENDPOINT`` at whichever tracing backend you use. To print spans straight to the terminal instead, with no backend at all, drop that variable and set ``OTEL_TRACES_EXPORTER=console`` in its place.
|
||||
|
||||
A few things catch people out the first time:
|
||||
|
||||
.. warning::
|
||||
``OTEL_TRACES_EXPORTER=console datasette mydb.db`` produces **nothing**. That environment variable is read by the OpenTelemetry SDK's auto-configuration, which only runs when the ``opentelemetry-instrument`` agent wraps the process. Datasette core installs no provider, so a plain ``datasette`` process emits nothing at all, whatever ``OTEL_`` variables are set.
|
||||
|
||||
- **Spans do not appear immediately.** The SDK's default ``BatchSpanProcessor`` flushes on a timer, every 5 seconds. Either wait, or stop the process - shutdown triggers a final flush - or set ``OTEL_BSP_SCHEDULE_DELAY=1000`` while you are experimenting. That last one is for demos, not for production.
|
||||
|
||||
- **Always set** ``OTEL_SERVICE_NAME``. Without it the SDK's default resource reports a ``service.name`` of ``unknown_service``, and your traces will be filed under that instead of under a name you can search for.
|
||||
|
||||
- **Setting** ``OTEL_LOGS_EXPORTER=none`` is worth doing unless your backend accepts logs too - ``opentelemetry-distro`` defaults every signal to OTLP, and a backend that does not take a signal will reject it noisily. Datasette emits no logs through OpenTelemetry; it does emit metrics (see :ref:`internals_telemetry_metrics`), so set ``OTEL_METRICS_EXPORTER=none`` only if your backend does not accept them.
|
||||
|
||||
Span reference
|
||||
--------------
|
||||
|
||||
Datasette emits six spans. One covers the HTTP request, and is the root everything else raised while serving that request hangs from. Four describe the database layer - one per query, one for the work that query does inside a SQL worker thread, and two more for the write queue. The sixth covers startup. Attribute names use the ``datasette.*`` prefix for Datasette-specific data, alongside standard OpenTelemetry attributes such as ``db.system``.
|
||||
|
||||
This reference is generated from ``datasette/telemetry_registry.py``, the single source of truth for every span and attribute Datasette emits. A conformance test makes real requests and compares what is actually emitted against that registry in both directions, so nothing here is hand-maintained and nothing can silently drift out of date.
|
||||
|
||||
Spans are ``SpanKind.INTERNAL`` unless a kind is listed below. Two are not: the request span is ``SERVER``, and ``db.query`` is ``CLIENT`` because it is the one span that represents a call to a database rather than Datasette's own work. Trace UIs use the kind to decide whether to render a span as an inbound request or as a database call. ``db.query``'s children stay ``INTERNAL`` because they are Datasette's decomposition of that one query - marking them ``CLIENT`` too would make a single query look like several database calls to anything counting by kind.
|
||||
|
||||
The request span's name is the only one that is not a fixed string - it is composed from the request, so the heading below shows the template rather than a literal you will see in a trace. A request to a table page produces a span named, in full::
|
||||
|
||||
GET /(?P<database>[^\/\.]+)/(?P<table>[^\/\.]+)(\.(?P<format>\w+))?$
|
||||
|
||||
That is the route's compiled regular expression, not a prettified ``/{database}/{table}`` template - see the ``http.route`` attribute below for why. Django's own instrumentation ships regex-flavoured routes for the same reason.
|
||||
|
||||
.. [[[cog
|
||||
from telemetry_doc import spans
|
||||
spans(cog)
|
||||
.. ]]]
|
||||
|
||||
``{http.request.method} {http.route}``
|
||||
One span per HTTP request, created by the outermost layer of the ASGI stack - so plugin ``asgi_wrapper()`` middleware, CSRF protection and every database span raised while serving the request all nest inside it. Without it each of those would be its own root trace. The span name is not a fixed string: it is the method followed by the matched route, and just the method for a request that matched no route. The span starts at the ASGI edge, before routing has happened, so it is named for the method there and renamed once the route is known. W3C ``traceparent`` and ``baggage`` headers are extracted using the global propagator, so a request arriving from an already-traced caller continues that trace; set ``OTEL_PROPAGATORS=none`` to turn that off, and strip those headers at your proxy if your instance is public.
|
||||
|
||||
Kind: ``SERVER``.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``http.request.method`` - The HTTP method, clamped to the nine methods RFC 9110 and RFC 5789 define. Anything else is reported as ``_OTHER``: the method is a client-controlled string, so echoing it back unbounded would be a cardinality hazard.
|
||||
- ``http.route`` *(optional)* - The route the request matched, as the compiled regular expression pattern Datasette routes with - for example ``/(?P<database>[^\/\.]+)/(?P<table>[^\/\.]+)(\.(?P<format>\w+))?$`` for a table page. It is deliberately the pattern rather than a prettified ``/{database}/{table}`` template: the route table is fixed when the app is built, so the pattern is exact, bounded and needs no parsing, whereas the transform into something prettier accretes edge cases. Unlike ``url.path`` this is low cardinality, so it is the attribute to group by. Omitted when no route matched - a 404 - which is also when the span name falls back to the bare method.
|
||||
- ``url.path`` - The path portion of the URL. The query string is deliberately **not** recorded, on this or any other span: Datasette puts user-supplied SQL in ``?sql=`` and canned query parameters in the query string, so exporting it by default would export exactly the data the rest of this instrumentation is careful with.
|
||||
- ``url.scheme`` - ``http`` or ``https``.
|
||||
- ``server.address`` *(optional)* - The ``Host`` header, verbatim - including any ``:port`` suffix, a deliberate deviation from semantic conventions' ``server.address`` / ``server.port`` split. Client-controlled, so treat it as untrusted input rather than as the identity of the server.
|
||||
- ``user_agent.original`` *(optional)* - The ``User-Agent`` header, verbatim. Omitted if the client sent none.
|
||||
- ``http.response.status_code`` *(optional)* - The status of the response, read from the ASGI ``http.response.start`` message rather than from a :ref:`internals_response` object - several views, including static files, file downloads and streaming CSV, send that message themselves and never build one. Omitted if the connection closed before anything was sent.
|
||||
- ``error.type`` *(optional)* - Set when the request failed: the exception class name if one escaped the application, otherwise the status code as a string for a 5xx response. A 4xx does **not** set this and does not set an error status - per semantic conventions a client error is not a server span's failure.
|
||||
- ``datasette.internal_client`` *(optional)* - ``True`` when the request was made in-process through ``datasette.client`` rather than arriving over the network. Such a sub-request runs the full ASGI stack, so it emits its own nested ``SERVER`` span inside the outer request's - filter on this attribute to keep kind-based dashboards from double-counting requests. Omitted for real inbound requests.
|
||||
|
||||
``db.query``
|
||||
A SQL operation issued by Datasette, covering the full round trip including any time spent queued for a thread. Callback-style calls - ``execute_fn()``, ``execute_write_fn()`` and ``execute_isolated_fn()`` - appear here too, distinguished by ``datasette.callback`` in place of ``db.query.text``.
|
||||
|
||||
Kind: ``CLIENT``.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``db.system`` - Always ``sqlite``.
|
||||
- ``db.namespace`` - Name of the database being queried.
|
||||
- ``db.query.text`` *(optional)* - The SQL, truncated to 2048 characters. Never the parameter values. Absent for a callback-style call (``execute_fn()`` and friends), where there is no SQL string to record - ``datasette.callback`` is set instead.
|
||||
- ``datasette.callback`` *(optional)* - The qualified name of the Python callable passed to ``execute_fn()``, ``execute_write_fn()`` or ``execute_isolated_fn()`` - for example ``TableInsertView.post.<locals>.insert_or_upsert_rows``. Set instead of ``db.query.text``, which does not exist for a callback: the SQL is whatever the function chooses to run. A lambda reports ``<lambda>``, which is why callers wanting a recognisable span should pass a named function. Bounded cardinality: the set of callables is fixed by the installed code, not by request input.
|
||||
- ``db.operation.name`` *(optional)* - The statement's leading keyword - ``SELECT``, ``INSERT``, ``CREATE``, and so on - matched against a small fixed allowlist. Omitted rather than set to an arbitrary value: the attribute must stay safe to use as a metric dimension, and echoing an unrecognised first token from user-supplied SQL would be an unbounded-cardinality hazard. Also omitted for ``execute_write_script()``, which runs multiple statements - per semantic conventions, the operation name should not be extracted from query text that can contain more than one operation. Note that a statement beginning with a CTE reports ``WITH``, not the operation inside it - a substantial share of Datasette's own reads take that form. Resolving it further would mean parsing.
|
||||
- ``db.collection.name`` *(optional)* - The primary table, set only where the view already knows it - the table and row pages. Omitted for arbitrary ``?sql=`` queries, where determining the table would mean parsing the query.
|
||||
- ``datasette.param_count`` *(optional)* - Number of bound parameters. Recorded instead of the values themselves.
|
||||
- ``datasette.param_sets`` *(optional)* - Number of parameter sets consumed by ``execute_write_many()``. Not a row count - ``executemany()`` returns no rows. The parameter values themselves are never recorded: that sequence can hold thousands of rows.
|
||||
- ``datasette.time_limit_ms`` *(optional)* - The :ref:`setting_sql_time_limit_ms` value this query ran under. Set on reads, which are the queries that time limit applies to.
|
||||
- ``datasette.rows_returned`` *(optional)* - Number of rows a read returned. Set on the read path only, and only when the read succeeded.
|
||||
- ``datasette.truncated`` *(optional)* - True if the result was cut short by :ref:`setting_max_returned_rows`.
|
||||
- ``datasette.interrupted`` *(optional)* - True if the query was cancelled for exceeding the time limit. The span status is also set to ``ERROR``, unless the caller asked for a budget shorter than :ref:`setting_sql_time_limit_ms` - as table counts, facet suggestion and autocomplete all do - in which case running out of time is an expected answer rather than a failure and the status is left unset.
|
||||
- ``datasette.sql_error_suppressed`` *(optional)* - True when the query failed but the caller passed ``log_sql_errors=False``, meaning it was probing and treats failure as an expected answer. Facet suggestion does this against every column.
|
||||
- ``datasette.executescript`` *(optional)* - True for ``execute_write_script()``, which runs multiple statements.
|
||||
- ``datasette.executemany`` *(optional)* - True for ``execute_write_many()``, which runs one statement against many parameter sets.
|
||||
|
||||
``db.query.execute``
|
||||
The read executing inside a SQL worker thread. Child of ``db.query``; the gap between the two is time spent waiting for a thread.
|
||||
|
||||
No attributes.
|
||||
|
||||
``db.write.queue_wait``
|
||||
Time a write spent waiting in its database's write queue before the write thread picked it up. Child of ``db.query`` for a ``block=True`` write, where the caller awaits the write and containment is accurate. For a ``block=False`` write the caller does not await it - the enqueueing request *caused* the write without *containing* it, and the write's spans can outlive the request's own - so this is a root span instead, carrying an OpenTelemetry link back to the enqueueing span rather than a parent. A link records causation without asserting containment, which is exactly the distinction here.
|
||||
|
||||
No attributes.
|
||||
|
||||
``db.write.execute``
|
||||
The write executing on the write thread. Child of ``db.query`` for a ``block=True`` write; for ``block=False`` a root span with a link back to the enqueueing span instead - see ``db.write.queue_wait`` above.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``datasette.isolated_connection`` - True if the write ran on its own connection rather than the shared write connection.
|
||||
- ``datasette.transaction`` - False for statements such as ``VACUUM`` that cannot run inside a transaction.
|
||||
|
||||
``datasette.startup``
|
||||
``invoke_startup()`` running: ``register_events``, ``register_actions``, ``register_column_types``, ``prepare_jinja2_environment``, internal-database schema catalog refresh (including the ``prepare_connection`` warm-up this triggers for each database touched for the first time), saved queries, column type config and the ``startup`` hook. Runs once per process, before any request exists, so without this span every child it creates would be its own orphan root trace. A connection warmed later - lazily, the first time a *request* touches a new database or thread - nests under that request's own span instead, not under this one, since this span has already ended by then.
|
||||
|
||||
No attributes.
|
||||
|
||||
.. [[[end]]]
|
||||
|
||||
.. _internals_telemetry_metrics:
|
||||
|
||||
Metric reference
|
||||
----------------
|
||||
|
||||
Spans describe events; metrics describe levels and rates. "Am I saturating my :ref:`setting_num_sql_threads` threads right now?" cannot be answered by any span, because it is a level sampled at collection time - and it is usually the first thing worth knowing about a busy Datasette, since ``num_sql_threads`` defaults to ``3``. Metrics also survive trace sampling: an operator keeping 1% of traces still gets 100% of every histogram and counter below.
|
||||
|
||||
As with spans, core emits these through the OpenTelemetry API only. Without a ``MeterProvider`` every instrument is a no-op, and the observable-gauge callbacks are never invoked at all, so an uninstrumented install pays nothing for them.
|
||||
|
||||
Every duration histogram is in **seconds**, with explicit bucket boundaries chosen for an in-process database - OpenTelemetry's default boundaries are tuned for milliseconds and would file every SQLite query into a single bucket, making quantile queries meaningless. The boundaries are listed with each histogram because a ``histogram_quantile()`` query is only as good as the buckets underneath it.
|
||||
|
||||
This reference is generated from ``datasette/telemetry_registry.py``, like the span reference above.
|
||||
|
||||
.. [[[cog
|
||||
from telemetry_doc import metrics
|
||||
metrics(cog)
|
||||
.. ]]]
|
||||
|
||||
``db.client.operation.duration``
|
||||
Histogram, unit ``s``. Duration of a SQL operation. The standard OpenTelemetry semantic convention metric, and the one that survives trace sampling. Callback-style calls (``execute_fn()`` and friends) are counted alongside the SQL-string methods.
|
||||
|
||||
Bucket boundaries: ``0.0001``, ``0.0005``, ``0.001``, ``0.005``, ``0.01``, ``0.05``, ``0.1``, ``0.5``, ``1``, ``5``, ``10``.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``db.system`` - Always ``sqlite``.
|
||||
- ``db.namespace`` - Name of the database being queried.
|
||||
- ``datasette.operation`` - Whether the operation was a read or a write. One of: ``read``, ``write``.
|
||||
- ``error.type`` *(optional)* - Set when the request failed: the exception class name if one escaped the application, otherwise the status code as a string for a 5xx response. A 4xx does **not** set this and does not set an error status - per semantic conventions a client error is not a server span's failure.
|
||||
|
||||
``datasette.write.queue_wait``
|
||||
Histogram, unit ``s``. Time each write waited in its database's write queue. The metric counterpart of the ``db.write.queue_wait`` span.
|
||||
|
||||
Bucket boundaries: ``0.0001``, ``0.0005``, ``0.001``, ``0.005``, ``0.01``, ``0.05``, ``0.1``, ``0.5``, ``1``, ``5``, ``10``.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``db.namespace`` - Name of the database being queried.
|
||||
|
||||
``datasette.sql.queries.interrupted``
|
||||
Counter, unit ``{query}``. Queries cancelled for exceeding :ref:`setting_sql_time_limit_ms`. Worth alerting on: a rising rate means the limit is too tight or a table has outgrown its queries. A caller that opted into a deliberately shorter budget - facet suggestion, for example - is not counted, for the same reason its timeout is not a span error.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``db.namespace`` - Name of the database being queried.
|
||||
|
||||
``datasette.sql.threads.limit``
|
||||
Observable gauge, unit ``{thread}``. Maximum concurrent read queries - the :ref:`setting_num_sql_threads` value. Not reported when ``num_sql_threads`` is ``0``, since then queries run on the event loop and there is no pool.
|
||||
|
||||
No attributes.
|
||||
|
||||
``datasette.sql.threads.queue_depth``
|
||||
Observable gauge, unit ``{query}``. Read queries waiting for a free thread. **This is the saturation signal** - sustained above zero means requests are queueing on ``num_sql_threads``.
|
||||
|
||||
No attributes.
|
||||
|
||||
``datasette.sql.queries.pending``
|
||||
Observable gauge, unit ``{query}``. Read queries submitted to the pool and not yet complete. Summed across databases and compared against the thread limit, this is pool utilisation.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``db.namespace`` - Name of the database being queried.
|
||||
|
||||
``datasette.write.queue_depth``
|
||||
Observable gauge, unit ``{write}``. Writes queued behind a database's single write thread. Backpressure that raising ``num_sql_threads`` cannot relieve. Not reported for a database that has never been written to.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``db.namespace`` - Name of the database being queried.
|
||||
|
||||
``datasette.connections.open``
|
||||
Observable gauge, unit ``{connection}``. Open SQLite file connections currently tracked for closing.
|
||||
|
||||
Attributes:
|
||||
|
||||
- ``db.namespace`` - Name of the database being queried.
|
||||
|
||||
.. [[[end]]]
|
||||
|
||||
Exemplars
|
||||
~~~~~~~~~
|
||||
|
||||
An OpenTelemetry `exemplar <https://opentelemetry.io/docs/specs/otel/metrics/data-model/#exemplars>`__ attaches a trace ID and span ID to one sample backing a histogram measurement. Where a spike in ``db.client.operation.duration`` alone tells you "queries were slow sometime in this minute", the exemplar attached to one of the samples in that spike gives you the trace ID of an actual slow query to open.
|
||||
|
||||
Datasette needs no configuration to produce these. Every histogram measurement on the query path is recorded while a span for that operation is active, and the OpenTelemetry SDK's default exemplar filter attaches the current trace ID and span ID to any measurement recorded inside a sampled span - this is SDK behaviour that Datasette's instrumentation does not need to opt into. Four queries of increasing cost, each inside its own span, produced one exemplar per query on ``db.client.operation.duration``:
|
||||
|
||||
.. code-block:: text
|
||||
|
||||
db.client.operation.duration count=4
|
||||
exemplars: 4
|
||||
value=0.001564s trace_id=ddfaf45fd4e14913497d7efeac95f381 span_id=fd5792bdbb01e533
|
||||
value=0.006320s trace_id=34aea775ade11a3c5f716695731000fe span_id=25ed9e29dd84dbee
|
||||
value=0.045253s trace_id=a65cb58d1460a179f0d04046ff51ed0d span_id=7f34d6378c85d062
|
||||
value=0.305240s trace_id=6089f4c515c221c0ca7bb53667b37ac8 span_id=0516f4a6641eaa0b
|
||||
|
||||
Exemplars are kept per histogram bucket - the SDK's default reservoir for an explicit-bucket histogram holds one exemplar per bucket - so the bucket boundaries above decide how many distinct traces a metric can point at. The same four queries, run against an earlier set of bucket boundaries under which every one of them fell into a single ``(0, 5]`` second bucket, produced one exemplar instead of four:
|
||||
|
||||
.. code-block:: text
|
||||
|
||||
db.client.operation.duration count=4
|
||||
exemplars: 1
|
||||
value=0.305349s trace_id=cd40f9af396ad1e1d71a8832d70ac84a span_id=8a8c288609731fea
|
||||
|
||||
Correcting the bucket boundaries had a second effect beyond fixing the quantiles: it also multiplied the number of traces reachable from this metric, one to four for this workload.
|
||||
|
||||
Which export path you use matters here. The OTLP exporter carries exemplars through unchanged. ``opentelemetry-exporter-prometheus`` (as of ``0.65b0``) does not: the string ``exemplar`` does not appear anywhere in its source, and an OpenMetrics scrape of the workload above through that exporter contained zero exemplar markers - even though ``prometheus-client`` (``0.26.0``), the library it depends on for OpenMetrics output, supports the syntax. If exemplars need to reach Prometheus, the path that works is an OTLP collector writing to Prometheus, not a plugin scraping through that exporter. On that path, the Prometheus server needs `--enable-feature=exemplar-storage <https://prometheus.io/docs/prometheus/latest/feature_flags/>`__, and the scrape itself must use the OpenMetrics exposition format - Prometheus's default text format has no syntax for exemplars at all. Grafana then needs the Prometheus data source's `exemplar configuration <https://grafana.com/docs/grafana/latest/datasources/prometheus/configure/>`__ (``exemplarTraceIdDestinations``) pointed at a tracing data source before it will draw an exemplar as a clickable point rather than an ordinary sample.
|
||||
|
||||
An exemplar can only exist for a trace that was sampled. The SDK's default exemplar filter only records one when the measurement happens inside a sampled span, and produces no exemplar at all rather than a link to a trace that was never kept. Verified: with the tracer provider's sampler set to ``ALWAYS_OFF``, the same four-query workload produced ``exemplars: 0`` on every data point. At 1% head sampling, 99% of measurements contribute no exemplar - but every exemplar you do get is guaranteed to resolve to a trace that exists.
|
||||
|
||||
.. _internals_telemetry_requests:
|
||||
|
||||
Requests and inbound trace context
|
||||
----------------------------------
|
||||
|
||||
Datasette creates the request span itself, at the outermost layer of the ASGI stack, so a trace is complete out of the box with no plugin and no extra instrumentation package. Everything raised while serving the request - plugin ``asgi_wrapper()`` middleware, CSRF protection, every database query - nests inside it.
|
||||
|
||||
**Inbound trace context is trusted by default.** W3C ``traceparent`` and ``baggage`` headers are extracted from every request using the global propagator, so a request arriving from an already-traced caller continues that trace instead of starting a new one. That is what every other framework instrumentation does - Flask, Django, FastAPI and ``opentelemetry-instrumentation-asgi`` all extract unconditionally - but on an instance open to the internet it means an arbitrary client can influence your traces:
|
||||
|
||||
- **Trace-ID pollution.** The client chooses the trace ID its request is filed under.
|
||||
- **Sampling control.** The SDK's default sampler is ``parentbased_always_on``, so under any parent-based sampler a client's sampled flag can force recording - a telemetry-cost denial of service - or suppress it.
|
||||
- **Baggage injection**, through the default composite propagator.
|
||||
|
||||
Because extraction goes through the *global* propagator there is no Datasette setting to configure, and the remedies are the standard OpenTelemetry ones:
|
||||
|
||||
- Strip ``traceparent``, ``tracestate`` and ``baggage`` at your reverse proxy, which is the right answer for a public instance fronted by one.
|
||||
- Set ``OTEL_PROPAGATORS=none`` to disable extraction entirely, or ``OTEL_PROPAGATORS=tracecontext`` to keep trace continuation and drop baggage.
|
||||
- Use a sampler that is not parent-based, which neutralises the sampling concern on its own.
|
||||
|
||||
**Installing an ASGI instrumentation as well is harmless.** If you wire up ``opentelemetry-instrumentation-asgi`` through an ``asgi_wrapper()`` plugin, its middleware lands *inside* Datasette's own, so its span becomes a redundant child ``SERVER`` span in the same trace. Nothing is re-orphaned. There is no setting to turn Datasette's request span off, because "turn it off" is already covered by installing no provider, or by ``OTEL_SDK_DISABLED=true``.
|
||||
|
||||
**Where** ``datasette.startup`` **lands depends on how you run Datasette.** ``datasette serve`` calls ``invoke_startup()`` before the server starts accepting connections, so the startup span is its own trace. An ASGI-hosted or programmatic deployment reaches startup lazily, on the first request, so there the startup span nests under that first request - which is honest, since it genuinely is that request's latency.
|
||||
|
||||
.. _internals_telemetry_privacy:
|
||||
|
||||
Privacy and safety
|
||||
------------------
|
||||
|
||||
Spans leave your infrastructure whenever you configure an exporter, so what goes into them is a security decision. Datasette's rules are:
|
||||
|
||||
- **SQL text is truncated to 2048 characters.** On a public instance the SQL is supplied by visitors and is unbounded in length, so ``db.query.text`` is cut off - with a ``…[truncated]`` marker - rather than allowed to set the size of a span.
|
||||
- **SQL parameter values are never recorded.** Only ``datasette.param_count``, a count. Parameter values are the part of a query most likely to hold something sensitive, and separating them from the SQL is the reason bound parameters exist.
|
||||
- **No actor identifiers are recorded.** No actor ID, no actor JSON, no client IP address. Nothing on a span identifies who made the request.
|
||||
- **Table names come only from an explicit** ``table=`` **argument.** ``db.collection.name`` is set by callers that already know which table they are working with, and is never derived from the SQL. Deriving it would mean parsing, and on an instance where visitors can create tables the set of possible values has no ceiling.
|
||||
- **The query string is never recorded.** There is no ``url.query`` attribute on the request span or on any other span. Datasette puts user-supplied SQL in ``?sql=`` and canned query parameters in the query string, so recording it by default would export exactly the class of data the rules above are careful with. Only ``url.path`` and ``http.route`` are recorded.
|
||||
|
||||
The SQL itself, though, *is* recorded, and on a public instance that means anything a visitor types into the query editor or passes as ``?sql=`` will be exported along with the span. That is the trade-off tracing a query engine makes.
|
||||
|
||||
.. _internals_telemetry_limitations:
|
||||
|
||||
Known limitations
|
||||
-----------------
|
||||
|
||||
- ``http.route`` **is a compiled regular expression, not a pretty route template.** See :ref:`internals_telemetry_requests` above for why.
|
||||
- **Inbound trace context is trusted by default**, which on a public instance means a client can influence your trace IDs, your sampling and your baggage. :ref:`internals_telemetry_requests` lists the remedies.
|
||||
- **Two plugin hooks run outside the** ``datasette.startup`` **span.** ``register_output_renderer`` is dispatched from ``Datasette.__init__()`` and ``asgi_wrapper`` from ``Datasette.app()``, both of which happen before ``invoke_startup()``. Datasette itself queries no database in either, so a default install emits nothing there - but a plugin that does will produce a root trace. Covering these would mean holding a span open across object construction, which is worse than the orphan.
|
||||
- ``db.operation.name`` **reports** ``WITH`` **for a statement that opens with a common table expression**, rather than the operation inside it, and a substantial share of Datasette's own reads take that form. The attribute is a leading-keyword match against a fixed allowlist, deliberately not a parse.
|
||||
- **Spans emitted before a provider is installed are not recorded.** If you are embedding Datasette in a host application, install your ``TracerProvider`` before serving traffic. This is ordinary OpenTelemetry behaviour rather than anything Datasette controls; nothing is permanently affected, those particular spans are simply dropped.
|
||||
|
||||
.. _internals_csrf:
|
||||
|
||||
CSRF protection
|
||||
|
|
|
|||
|
|
@ -1,224 +0,0 @@
|
|||
.. _plugin_telemetry:
|
||||
|
||||
Telemetry for plugin authors
|
||||
============================
|
||||
|
||||
Datasette core emits OpenTelemetry spans and metrics for the work it does itself - see :ref:`internals_telemetry` for what those are and how an operator turns them on. This page is about the other half: instrumenting the work **your plugin** does, so that a plugin's queries, background jobs and custom operations show up in the same traces and the same metrics pipeline, using the same conventions.
|
||||
|
||||
Everything here follows one rule inherited from core: **depend on** ``opentelemetry-api`` **only, and never install a provider**. With no SDK installed every span and instrument your plugin creates is a free no-op; whoever runs Datasette decides whether telemetry is collected, sampled or exported. A plugin that installs a ``TracerProvider`` or configures an exporter is making an operator's decision for them.
|
||||
|
||||
.. _plugin_telemetry_scope:
|
||||
|
||||
Use your own instrumentation scope
|
||||
----------------------------------
|
||||
|
||||
Never emit through core's tracer or meter. Your plugin's scope name is the machine-readable claim about *who emitted a signal*, and consumers filter on it:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
from opentelemetry import metrics, trace
|
||||
|
||||
from my_plugin import __version__
|
||||
|
||||
tracer = trace.get_tracer("my_plugin", __version__)
|
||||
meter = metrics.get_meter("my_plugin", __version__)
|
||||
|
||||
If every attribute you emit follows current semantic conventions you can also pass ``schema_url=``; ``datasette.telemetry.SCHEMA_URL`` is the version core's own spellings track, with a comment explaining how to choose one. When in doubt, omit it - a wrong schema URL is worse than none.
|
||||
|
||||
Two naming rules keep the ecosystem's signals tellable-apart:
|
||||
|
||||
- **Scope**: use your plugin's *import package* name - ``my_plugin``, underscores and all. Consumers filter on the scope, and one spelling convention means they can guess it.
|
||||
- **Signal prefix**: name spans, metrics and custom attributes under a prefix you own - your package name (``my_plugin.*``) or a short product name (``paper.*``). **Never a bare** ``datasette.*`` **prefix**: that namespace belongs to core, an operator could no longer tell core signals from plugin signals, and a future core signal could collide with yours.
|
||||
|
||||
Reuse core's shared attribute spellings where they mean the same thing - ``db.namespace`` for a database name, ``error.type`` for a failure class - rather than minting parallel ones. If a span family in your registry shares a prefix with another entry, exact names always win over prefix matches, but two overlapping ``prefix=True`` entries resolve to whichever is listed first - avoid overlapping families rather than relying on order.
|
||||
|
||||
.. _plugin_telemetry_registry:
|
||||
|
||||
Declare a registry
|
||||
------------------
|
||||
|
||||
Core keeps a single source of truth for every signal it emits in ``datasette/telemetry_registry.py``, and the classes it uses are public API. They subclass ``str``, so a registry entry *is* the name you pass to OpenTelemetry - no parallel constants to keep in step:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
from opentelemetry.trace import SpanKind
|
||||
|
||||
from datasette.telemetry_registry import Attribute, MetricName, SpanName
|
||||
|
||||
OUTCOME = Attribute(
|
||||
"my_plugin.outcome",
|
||||
"How the job ended.",
|
||||
values={"ok", "error", "skipped"},
|
||||
)
|
||||
JOB_NAME = Attribute("my_plugin.job", "The registered job name.")
|
||||
|
||||
JOB_RUN = SpanName(
|
||||
"my_plugin.job.run",
|
||||
"One execution of a scheduled job.",
|
||||
(OUTCOME, JOB_NAME),
|
||||
)
|
||||
|
||||
# A span family with a variable suffix - emitted as "my_plugin.chat gpt-5"
|
||||
CHAT = SpanName(
|
||||
"my_plugin.chat ",
|
||||
"One model call, named ``my_plugin.chat {model}``.",
|
||||
prefix=True,
|
||||
)
|
||||
|
||||
SPANS = (JOB_RUN, CHAT)
|
||||
|
||||
JOB_DURATION = MetricName(
|
||||
"my_plugin.job.duration",
|
||||
"Histogram",
|
||||
"s",
|
||||
"How long each job took.",
|
||||
(JOB_NAME, OUTCOME),
|
||||
buckets=(0.01, 0.1, 1, 10, 60, 600, 3600),
|
||||
)
|
||||
|
||||
Three details that matter:
|
||||
|
||||
- ``values=`` declares a **closed enum**. The conformance helpers (below) assert every emitted value is a member, which is what makes an attribute safe to use as a metric dimension - a metric series is keyed by its attribute values, so an open value set on a metric is an unbounded-cardinality hazard.
|
||||
- ``prefix=True`` registers a span *family* whose emitted names share a fixed prefix; ``datasette.telemetry_registry.span_for()`` matches them by prefix, exact names first.
|
||||
- Declare explicit histogram ``buckets=`` scaled to *your* domain. Core's SQLite-scale boundaries are importable as ``datasette.telemetry_registry.DURATION_BUCKETS`` (0.0001s to 10s) - use them if you are timing SQLite work so dashboards align, and define your own otherwise (a job scheduler wants buckets out to an hour; the SDK's defaults will put all your measurements in one bucket either way).
|
||||
|
||||
.. _plugin_telemetry_privacy:
|
||||
|
||||
Privacy and cardinality rules
|
||||
-----------------------------
|
||||
|
||||
Core's instrumentation records **no data users put into Datasette and no identifier that ties a signal to a person** - no parameter values, no query strings, no actor identifiers, no IP addresses. Hold your plugin to the same bar:
|
||||
|
||||
- Attribute values should be closed enums, booleans, counts and durations. Anything echoed from user input - a name, a URL, a token, free text - does not belong on a span, and *especially* not on a metric.
|
||||
- If you time user-influenced SQL, follow core: record the SQL via ``datasette.telemetry.sql_attribute()`` (truncated, never parameters) on spans only.
|
||||
- When a value is interesting but unbounded, record a bounded proxy instead: a count, a byte size, a truncation flag, or the enum outcome.
|
||||
|
||||
These rules are enforceable: see ``assert_no_forbidden_values()`` in :ref:`plugin_telemetry_testing`.
|
||||
|
||||
.. _plugin_telemetry_callbacks:
|
||||
|
||||
Your database work is already traced
|
||||
------------------------------------
|
||||
|
||||
Every call your plugin makes through :ref:`db.execute() <database_execute>`, :ref:`db.execute_fn() <database_execute_fn>`, :ref:`db.execute_write() <database_execute_write>` and :ref:`db.execute_write_fn() <database_execute_write_fn>` already emits core's ``db.query`` spans and is counted in the ``db.client.operation.duration`` histogram. Two consequences:
|
||||
|
||||
- **Pass named callables**, not lambdas: the span for a callback-style call is identified by ``datasette.callback``, the callable's qualified name, and a lambda reports ``<lambda>``.
|
||||
- If you also wrap those calls in your own span or histogram, you are creating a *second* series in *your* scope - that is fine and sometimes right (yours can carry plugin-level attributes core cannot know), but it is a deliberate two-series design, not a substitute for core's.
|
||||
|
||||
.. _plugin_telemetry_request_span:
|
||||
|
||||
Enriching the request span
|
||||
--------------------------
|
||||
|
||||
Inside a view or ASGI middleware, ``datasette.telemetry.request_span(scope)`` returns the recording ``SERVER`` span for the current request, or ``None`` when nothing is recording - which is also your signal to skip any work done only to compute attributes:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
from datasette.telemetry import request_span
|
||||
|
||||
async def my_view(request):
|
||||
span = request_span(request.scope)
|
||||
if span is not None:
|
||||
span.set_attribute("my_plugin.cache", "hit")
|
||||
...
|
||||
|
||||
.. _plugin_telemetry_background:
|
||||
|
||||
Background work: roots with links
|
||||
---------------------------------
|
||||
|
||||
A background job, a scheduled task or a queue consumer must **not** parent its spans to the request that caused it - by the time the work runs, that request span has usually ended, and a child outliving its closed parent renders badly in every major trace UI. The correct shape, the one core itself uses for ``execute_write(block=False)``, is a **root span carrying a link** to the causing span:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
from datasette.telemetry import linked_root_span_kwargs
|
||||
|
||||
# Capture at scheduling time, while the causing span is current:
|
||||
kwargs = linked_root_span_kwargs()
|
||||
|
||||
# Later, wherever the work actually runs:
|
||||
with tracer.start_as_current_span("my_plugin.job.run", **kwargs) as span:
|
||||
span.set_attribute(OUTCOME, "ok")
|
||||
|
||||
For a periodic loop (a health check, a scheduler tick), the convention is one root span **per tick**, always emitted - including no-op ticks, with an outcome attribute saying so - plus a tick counter metric. Suppressing quiet ticks seems tidy but destroys the signal operators actually want: "is the loop still running?". Pair the spans with a gauge for the loop's staleness if the interval is long.
|
||||
|
||||
Two propagation facts worth knowing (details in ``datasette/telemetry.py``):
|
||||
|
||||
- Core's ``tracer`` and yours are proxies. A ``ProxyTracer`` permanently caches the first concrete tracer it resolves *after* a provider exists, so in embedded deployments the provider must be installed before the first span - importing the module is fine, starting spans is not. Meters forward retroactively; tracers do not.
|
||||
- ``asyncio.create_task`` copies the ambient context, so a long-running task created during a request will silently parent to that request's span - exactly the bug ``linked_root_span_kwargs()`` exists to avoid.
|
||||
|
||||
.. _plugin_telemetry_gauges:
|
||||
|
||||
Observable gauges
|
||||
-----------------
|
||||
|
||||
For a *level* - how many streams are open, how deep is a queue - register an observable gauge whose callback the SDK invokes on its own collection cycle. Three disciplines, all inherited from how core implements its pool gauges in ``datasette/telemetry.py``:
|
||||
|
||||
- Hold live objects **weakly** (a ``weakref.WeakSet`` guarded by a lock), so instrumenting an object never keeps it alive, and unregister on close.
|
||||
- The callback runs on the SDK's **collection thread**: never take a lock the request path holds, never await, never do I/O. Read cached state and yield ``Observation`` values; if freshness matters, refresh the cache from your own code and expose its staleness as another gauge.
|
||||
- With no provider installed the callback is **never invoked at all**, so gauges are free by default.
|
||||
|
||||
.. _plugin_telemetry_testing:
|
||||
|
||||
Testing your instrumentation
|
||||
----------------------------
|
||||
|
||||
``datasette.telemetry_testing`` ships the same fixtures and checks core's own suite uses. In your ``conftest.py``:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
from datasette.telemetry_testing import ( # noqa: F401
|
||||
otel_metrics,
|
||||
otel_meter_provider,
|
||||
otel_provider,
|
||||
otel_reset,
|
||||
otel_spans,
|
||||
)
|
||||
|
||||
``otel_provider`` and ``otel_meter_provider`` are session-scoped and autouse - they install a real SDK provider (in-memory, synchronous export) once per process, and do nothing when the SDK is not installed, so add ``opentelemetry-sdk`` to your test dependencies only. If a different provider was installed first (an embedding app, ``opentelemetry-instrument``), the fixtures detect that the install did not take and skip with a clear message rather than asserting against an exporter wired to nothing. ``otel_reset`` is autouse too: it drains the exporter and reader after every test, so a large suite does not accumulate recorded spans for its whole lifetime. Tests then take ``otel_spans`` (an ``InMemorySpanExporter``) or ``otel_metrics`` (a collector with ``collect()`` / ``point()`` helpers).
|
||||
|
||||
Wire your registry to reality with the conformance helpers - the two directions catch instrumentation added without documentation and documentation describing signals that no longer exist:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
from datasette.telemetry_testing import (
|
||||
assert_metrics_conform,
|
||||
assert_metrics_covered,
|
||||
assert_package_never_imports_sdk,
|
||||
assert_spans_covered,
|
||||
assert_spans_conform,
|
||||
)
|
||||
|
||||
from my_plugin.telemetry import METRICS, SPANS
|
||||
|
||||
|
||||
def test_conformance(otel_spans, otel_metrics):
|
||||
run_a_workload_that_exercises_everything()
|
||||
finished = otel_spans.get_finished_spans()
|
||||
# Everything emitted is registered (and enum values are legal):
|
||||
assert_spans_conform(SPANS, finished, scope_name="my_plugin")
|
||||
# Everything registered was emitted:
|
||||
assert_spans_covered(SPANS, finished, scope_name="my_plugin")
|
||||
# Same two directions for metrics - one collect() after the workload:
|
||||
otel_metrics.collect()
|
||||
assert_metrics_conform(METRICS, otel_metrics, scope_name="my_plugin")
|
||||
assert_metrics_covered(METRICS, otel_metrics, scope_name="my_plugin")
|
||||
|
||||
|
||||
def test_api_only_dependency():
|
||||
assert_package_never_imports_sdk("my_plugin")
|
||||
|
||||
Always pass ``scope_name`` - the exporter and reader also hold core's signals, and your registry should only be judged against your own.
|
||||
|
||||
The metric helpers check more than names: ``assert_metrics_conform`` asserts each instrument was created as the **kind** and **unit** its registry entry declares (the registry entry and the ``meter.create_*()`` call are separate statements, and a dashboard built on the registry's word breaks silently if they drift), and that every value on a ``values=`` enum attribute is a member - which is what makes a metric dimension *provably* bounded rather than bounded by intent. Both ``*_covered`` helpers exempt attributes marked ``optional=True`` (an ``error.type`` only present on failures should not force your workload to manufacture errors - pin those with targeted tests instead), and the metrics reader uses delta temporality, so run one broad workload followed by a single ``collect()``.
|
||||
|
||||
Finally, enforce the privacy rules with ``assert_no_forbidden_values()``: plant sentinel values in your workload - a fake email your fixtures log in with, a token, a username - and assert they never appear in any span name, attribute, event, status description or metric attribute. Leave ``scope_name`` unset for this one: a secret leaking through *core's* signals (SQL text, say) is still a leak. Schedule the test that calls ``assert_package_never_imports_sdk()`` early in your suite - see its docstring for the macOS threading hazard.
|
||||
|
||||
.. _plugin_telemetry_caveats:
|
||||
|
||||
Known caveats
|
||||
-------------
|
||||
|
||||
- **Streaming responses hold the request span open.** Core's request span ends when the response body finishes, so for an SSE or long-streaming route its duration is the connection lifetime. If you need per-message timing on a stream, emit your own child spans or span events per message, and use gauges for concurrent-stream counts.
|
||||
- **A plugin timing core's work double-measures by design.** See :ref:`plugin_telemetry_callbacks` above.
|
||||
- ``datasette.client`` requests made from inside a request produce a nested ``SERVER`` span. Those spans carry ``datasette.internal_client: true`` - filter on it to keep kind-based dashboards from double-counting requests.
|
||||
|
|
@ -1,53 +0,0 @@
|
|||
"""
|
||||
Render the span reference in ``internals.rst`` from
|
||||
``datasette/telemetry_registry.py``.
|
||||
|
||||
Driven by cog, and ``cog --check docs/*.rst`` runs in CI - so adding a span
|
||||
without documenting it, or documenting one that no longer exists, is a build
|
||||
failure rather than something a reader discovers later.
|
||||
"""
|
||||
|
||||
|
||||
def _attribute_lines(cog, attributes):
|
||||
if not attributes:
|
||||
cog.out(" No attributes.\n\n")
|
||||
return
|
||||
cog.out(" Attributes:\n\n")
|
||||
for attribute in attributes:
|
||||
suffix = " *(optional)*" if attribute.optional else ""
|
||||
line = f" - ``{attribute}``{suffix} - {attribute.description}"
|
||||
if attribute.values is not None:
|
||||
rendered = ", ".join(f"``{value}``" for value in sorted(attribute.values))
|
||||
line += f" One of: {rendered}."
|
||||
cog.out(line + "\n")
|
||||
cog.out("\n")
|
||||
|
||||
|
||||
def spans(cog):
|
||||
from opentelemetry.trace import SpanKind
|
||||
|
||||
from datasette.telemetry_registry import SPANS
|
||||
|
||||
cog.out("\n")
|
||||
for span in SPANS:
|
||||
cog.out(f"``{span}``\n")
|
||||
cog.out(f" {span.description}\n\n")
|
||||
# INTERNAL is the default and the overwhelming majority of spans -
|
||||
# printing it on every one would be noise. Only the exceptional case,
|
||||
# a real database call, is worth calling out.
|
||||
if span.kind != SpanKind.INTERNAL:
|
||||
cog.out(f" Kind: ``{span.kind.name}``.\n\n")
|
||||
_attribute_lines(cog, span.attributes)
|
||||
|
||||
|
||||
def metrics(cog):
|
||||
from datasette.telemetry_registry import METRICS
|
||||
|
||||
cog.out("\n")
|
||||
for metric in METRICS:
|
||||
cog.out(f"``{metric}``\n")
|
||||
cog.out(f" {metric.kind}, unit ``{metric.unit}``. {metric.description}\n\n")
|
||||
if metric.buckets:
|
||||
boundaries = ", ".join(f"``{boundary}``" for boundary in metric.buckets)
|
||||
cog.out(f" Bucket boundaries: {boundaries}.\n\n")
|
||||
_attribute_lines(cog, metric.attributes)
|
||||
|
|
@ -40,7 +40,6 @@ dependencies = [
|
|||
"setuptools",
|
||||
"pip",
|
||||
"pydantic>=2",
|
||||
"opentelemetry-api>=1.37",
|
||||
]
|
||||
|
||||
[project.urls]
|
||||
|
|
@ -71,7 +70,6 @@ dev = [
|
|||
"cogapp>=3.3.0",
|
||||
"multipart-form-data-conformance==0.1a0",
|
||||
"ruff>=0.16.0",
|
||||
"opentelemetry-sdk>=1.37",
|
||||
# docs
|
||||
"Sphinx==7.4.7",
|
||||
"furo==2025.9.25",
|
||||
|
|
|
|||
|
|
@ -58,18 +58,6 @@ def find_free_port():
|
|||
return sock.getsockname()[1]
|
||||
|
||||
|
||||
# The otel fixtures moved to datasette.telemetry_testing, which is public
|
||||
# plugin API - core's suite consumes it exactly the way a plugin's would.
|
||||
from datasette.telemetry_testing import ( # noqa: F401, E402
|
||||
MetricsCollector,
|
||||
otel_metrics,
|
||||
otel_meter_provider,
|
||||
otel_provider,
|
||||
otel_reset,
|
||||
otel_spans,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def bare_ds():
|
||||
"""
|
||||
|
|
@ -180,14 +168,6 @@ def pytest_collection_modifyitems(config, items):
|
|||
move_to_front(items, "test_spatialite_error_if_attempt_to_open_spatialite")
|
||||
move_to_front(items, "test_package")
|
||||
move_to_front(items, "test_package_with_port")
|
||||
# Same reason: this one shells out to a fresh interpreter. Late in a serial
|
||||
# run the pytest process holds enough threads that the fork half of
|
||||
# subprocess' fork+exec crashes the interpreter on macOS/CPython 3.13
|
||||
# (SIGSEGV/SIGBUS inside _execute_child). Reproduces with any subprocess
|
||||
# call placed there, on an unmodified tree - running it first avoids it.
|
||||
move_to_front(items, "test_datasette_package_never_imports_the_sdk")
|
||||
move_to_front(items, "test_kit_module_itself_never_imports_the_sdk")
|
||||
move_to_front(items, "test_no_provider_takes_the_fast_path")
|
||||
|
||||
|
||||
def move_to_front(items, test_name):
|
||||
|
|
|
|||
|
|
@ -1,889 +0,0 @@
|
|||
"""
|
||||
The HTTP request span.
|
||||
|
||||
`tests/test_telemetry_registry.py` already pins the span's name shape, kind
|
||||
and attribute keys against literals, so this file deliberately does not
|
||||
repeat that. What it covers is the properties of the middleware and of the
|
||||
router's `http.route` enrichment that the registry conformance test
|
||||
structurally cannot see:
|
||||
|
||||
- **where the middleware sits.** Outermost is the entire point - moving it
|
||||
inside the plugin `asgi_wrapper()` loop leaves plugin middleware creating
|
||||
orphan root traces, which is the problem this span exists to fix, and every
|
||||
attribute assertion still passes.
|
||||
- **which span the route lands on**, which only diverges once something else
|
||||
has made a span current.
|
||||
- **method clamping**, which a workload of ordinary GETs can never exercise.
|
||||
- **the query string never being recorded**, which only fails if a request
|
||||
actually carries one.
|
||||
- **the span outliving a streamed response body**, which only a paging export
|
||||
can distinguish from ending far too early.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import itertools
|
||||
import json
|
||||
import subprocess
|
||||
import sys
|
||||
import textwrap
|
||||
import time
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
|
||||
pytest.importorskip("opentelemetry.sdk")
|
||||
|
||||
from opentelemetry.trace import (
|
||||
NonRecordingSpan,
|
||||
SpanContext,
|
||||
SpanKind,
|
||||
StatusCode,
|
||||
TraceFlags,
|
||||
)
|
||||
|
||||
from datasette import hookimpl
|
||||
from datasette.app import Datasette
|
||||
from datasette.telemetry import (
|
||||
REQUEST_SPAN_SCOPE_KEY,
|
||||
TelemetryMiddleware,
|
||||
request_span,
|
||||
tracer,
|
||||
)
|
||||
from datasette.utils import resolve_routes
|
||||
|
||||
# Named in-memory databases are shared-cache: two Datasette instances given
|
||||
# the same name share one SQLite database and the second `create table`
|
||||
# fails.
|
||||
_names = itertools.count()
|
||||
|
||||
|
||||
PLUGIN_MIDDLEWARE_SPAN = "test.plugin.middleware"
|
||||
|
||||
|
||||
class _MiddlewarePlugin:
|
||||
"A plugin asgi_wrapper() that creates a span, standing in for a real one."
|
||||
|
||||
__name__ = "HttpSpanMiddlewarePlugin"
|
||||
|
||||
@hookimpl
|
||||
def asgi_wrapper(self, datasette):
|
||||
def wrap(app):
|
||||
async def wrapped(scope, receive, send):
|
||||
with tracer.start_as_current_span(PLUGIN_MIDDLEWARE_SPAN):
|
||||
await app(scope, receive, send)
|
||||
|
||||
return wrapped
|
||||
|
||||
return wrap
|
||||
|
||||
|
||||
class _RaisingMiddlewarePlugin:
|
||||
"""
|
||||
A plugin asgi_wrapper() that raises.
|
||||
|
||||
`route_path` converts almost every exception into a 500 itself, so an
|
||||
exception escaping into the request span is only reachable from *outside*
|
||||
the router - a plugin wrapper, or a failure inside the 500 handler.
|
||||
"""
|
||||
|
||||
__name__ = "HttpSpanRaisingMiddlewarePlugin"
|
||||
|
||||
def __init__(self, call_app_first):
|
||||
self.call_app_first = call_app_first
|
||||
|
||||
@hookimpl
|
||||
def asgi_wrapper(self, datasette):
|
||||
call_app_first = self.call_app_first
|
||||
|
||||
def wrap(app):
|
||||
async def wrapped(scope, receive, send):
|
||||
if call_app_first:
|
||||
await app(scope, receive, send)
|
||||
raise RuntimeError("wrapper exploded")
|
||||
|
||||
return wrapped
|
||||
|
||||
return wrap
|
||||
|
||||
|
||||
class _BoomPlugin:
|
||||
"A route that raises, which route_path turns into a 500."
|
||||
|
||||
__name__ = "HttpSpanBoomPlugin"
|
||||
|
||||
@hookimpl
|
||||
def register_routes(self):
|
||||
return [(r"^/-/http-span-boom$", lambda: 1 / 0)]
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def ds():
|
||||
name = f"httpspan{next(_names)}"
|
||||
instance = Datasette(memory=True)
|
||||
instance.add_memory_database(name)
|
||||
await instance.invoke_startup()
|
||||
await instance.get_database(name).execute_write(
|
||||
"create table t (id integer primary key, v text)"
|
||||
)
|
||||
instance.db_name = name
|
||||
try:
|
||||
yield instance
|
||||
finally:
|
||||
instance.close()
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def ds_paging():
|
||||
"""
|
||||
An instance whose table is bigger than `max_returned_rows`.
|
||||
|
||||
That is what makes `?_stream=1` genuinely page: `stream_csv` loops calling
|
||||
`fetch_data` for each page *inside* the response body send, so the trace
|
||||
contains `db.query` spans that start after the response has begun. On a
|
||||
table that fits in one page every query finishes before the body starts
|
||||
and the span-covers-the-body assertion cannot fail.
|
||||
"""
|
||||
name = f"httpspanpaging{next(_names)}"
|
||||
# Both settings matter. `?_stream=1` forces `_size=max`, which is
|
||||
# `max_returned_rows` - so lowering only that gives one page of five rows
|
||||
# and no `next` token, and the export never loops.
|
||||
instance = Datasette(
|
||||
memory=True, settings={"max_returned_rows": 5, "default_page_size": 3}
|
||||
)
|
||||
instance.add_memory_database(name)
|
||||
await instance.invoke_startup()
|
||||
db = instance.get_database(name)
|
||||
await db.execute_write("create table t (id integer primary key, v text)")
|
||||
await db.execute_write_many(
|
||||
"insert into t (id, v) values (?, ?)", [[i, f"v{i}"] for i in range(40)]
|
||||
)
|
||||
instance.db_name = name
|
||||
try:
|
||||
yield instance
|
||||
finally:
|
||||
instance.close()
|
||||
|
||||
|
||||
def _server_spans(otel_spans):
|
||||
return [
|
||||
span for span in otel_spans.get_finished_spans() if span.kind is SpanKind.SERVER
|
||||
]
|
||||
|
||||
|
||||
def _route_for(ds, path):
|
||||
"The compiled pattern Datasette's own router resolves `path` to."
|
||||
match, _view = resolve_routes(ds._routes(), path)
|
||||
assert match is not None, f"{path} matches no route"
|
||||
return match.re.pattern
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_plugin_asgi_wrapper_middleware_runs_inside_the_request_span(
|
||||
ds, otel_spans
|
||||
):
|
||||
"""
|
||||
The placement check.
|
||||
|
||||
A span created by a plugin `asgi_wrapper()` must be a *child* of the
|
||||
request span. If the middleware is mounted anywhere inside the plugin
|
||||
loop the two swap places - the plugin's span becomes the root and the
|
||||
request span its child - which is exactly the orphaning this is meant to
|
||||
prevent, and which no attribute assertion notices.
|
||||
"""
|
||||
ds.pm.register(_MiddlewarePlugin(), name="httpspan-middleware")
|
||||
try:
|
||||
otel_spans.clear()
|
||||
response = await ds.client.get(f"/{ds.db_name}/t")
|
||||
assert response.status_code == 200
|
||||
finally:
|
||||
ds.pm.unregister(name="httpspan-middleware")
|
||||
|
||||
spans = otel_spans.get_finished_spans()
|
||||
server = [span for span in spans if span.kind is SpanKind.SERVER]
|
||||
assert len(server) == 1, "expected exactly one SERVER span per request"
|
||||
server_span = server[0]
|
||||
assert server_span.parent is None, "the request span should be the trace root"
|
||||
|
||||
plugin_spans = [span for span in spans if span.name == PLUGIN_MIDDLEWARE_SPAN]
|
||||
assert len(plugin_spans) == 1
|
||||
assert plugin_spans[0].parent is not None
|
||||
assert plugin_spans[0].parent.span_id == server_span.context.span_id
|
||||
assert plugin_spans[0].context.trace_id == server_span.context.trace_id
|
||||
|
||||
# And the database work is in the same trace, not off on its own.
|
||||
queries = [span for span in spans if span.name == "db.query"]
|
||||
assert queries, "a table page should have issued at least one query"
|
||||
for query in queries:
|
||||
assert query.context.trace_id == server_span.context.trace_id
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_unrecognised_method_is_clamped(ds, otel_spans):
|
||||
"""
|
||||
Anyone can send `FROB / HTTP/1.1`. An unclamped method is an unbounded
|
||||
dimension a client controls, so semantic conventions map anything off the
|
||||
known list to `_OTHER`.
|
||||
|
||||
The span name is checked too, and it is the reason the router clamps the
|
||||
method a second time when it renames the span: the middleware's clamping
|
||||
protects the attribute, but the name is rebuilt from `request.method` in
|
||||
`route_path`, which is the raw client string. An unclamped rename would
|
||||
put attacker-supplied text straight back into the span name.
|
||||
"""
|
||||
otel_spans.clear()
|
||||
await ds.client.request("FROB", f"/{ds.db_name}/t")
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["http.request.method"] == "_OTHER"
|
||||
assert server[0].name == f"_OTHER {server[0].attributes['http.route']}"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_known_method_is_not_clamped(ds, otel_spans):
|
||||
"The other half of clamping: a real method must survive it verbatim."
|
||||
otel_spans.clear()
|
||||
await ds.client.get(f"/{ds.db_name}/t")
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["http.request.method"] == "GET"
|
||||
assert server[0].name == f"GET {server[0].attributes['http.route']}"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_the_query_string_is_never_recorded(ds, otel_spans):
|
||||
"""
|
||||
Datasette puts user-supplied SQL in `?sql=` and canned query parameters in
|
||||
the query string, so no span may carry it. Asserting on the absence of a
|
||||
`url.query` key alone would not catch it arriving under some other name,
|
||||
so this searches every attribute value of every span for the marker.
|
||||
"""
|
||||
marker = "canary-9f2b1c"
|
||||
otel_spans.clear()
|
||||
await ds.client.get(f"/{ds.db_name}/t?_facet=v&_nosuch={marker}")
|
||||
spans = otel_spans.get_finished_spans()
|
||||
assert _server_spans(otel_spans), "no request span was emitted"
|
||||
leaked = [
|
||||
f"{span.name} -> {key}={value!r}"
|
||||
for span in spans
|
||||
for key, value in (span.attributes or {}).items()
|
||||
if marker in str(value) or key == "url.query"
|
||||
]
|
||||
assert not leaked, "the query string reached a span attribute: " + ", ".join(leaked)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_url_path_is_recorded_without_the_query_string(ds, otel_spans):
|
||||
otel_spans.clear()
|
||||
await ds.client.get(f"/{ds.db_name}/t?_facet=v")
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["url.path"] == f"/{ds.db_name}/t"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_escaping_exception_sets_error_type_and_reraises(ds, otel_spans):
|
||||
"""
|
||||
An exception that gets past `route_path` must be recorded, not swallowed.
|
||||
|
||||
No response ever started, so there is no status code to record either.
|
||||
"""
|
||||
ds.pm.register(
|
||||
_RaisingMiddlewarePlugin(call_app_first=False), name="httpspan-raiser"
|
||||
)
|
||||
try:
|
||||
otel_spans.clear()
|
||||
with pytest.raises(RuntimeError):
|
||||
await ds.client.get(f"/{ds.db_name}/t")
|
||||
finally:
|
||||
ds.pm.unregister(name="httpspan-raiser")
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["error.type"] == "RuntimeError"
|
||||
assert "http.response.status_code" not in server[0].attributes
|
||||
assert server[0].status.status_code is StatusCode.ERROR
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_escaping_exception_beats_the_status_code_for_error_type(
|
||||
ds, otel_spans
|
||||
):
|
||||
"""
|
||||
Both paths can fire on one request: a 500 response is sent and *then*
|
||||
something raises on the way out. The `finally` block runs while the
|
||||
exception is propagating, so without the guard it would overwrite the
|
||||
exception's class name with the string "500" - strictly less information
|
||||
about what actually went wrong.
|
||||
"""
|
||||
ds.pm.register(_BoomPlugin(), name="httpspan-boom")
|
||||
ds.pm.register(
|
||||
_RaisingMiddlewarePlugin(call_app_first=True), name="httpspan-raiser"
|
||||
)
|
||||
try:
|
||||
otel_spans.clear()
|
||||
with pytest.raises(RuntimeError):
|
||||
await ds.client.get("/-/http-span-boom")
|
||||
finally:
|
||||
ds.pm.unregister(name="httpspan-raiser")
|
||||
ds.pm.unregister(name="httpspan-boom")
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
# The 500 really was sent, so the status is still recorded ...
|
||||
assert server[0].attributes["http.response.status_code"] == 500
|
||||
# ... but error.type names the exception, not the status.
|
||||
assert server[0].attributes["error.type"] == "RuntimeError"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_a_404_is_not_an_error(ds, otel_spans):
|
||||
"""
|
||||
Per semantic conventions a 4xx is the client's mistake, not the server's,
|
||||
so a SERVER span must record the status and leave both its own status and
|
||||
`error.type` alone. Datasette 404s are routine - every missing table, and
|
||||
every bot probing for /wp-login.php - so treating them as errors would
|
||||
drown a real 500 in noise.
|
||||
|
||||
Note this 404 *does* match a route: `/no-such-database-at-all` matches the
|
||||
database pattern and the view then raises `NotFound`. Most Datasette 404s
|
||||
are that shape rather than the unrouted one below.
|
||||
"""
|
||||
otel_spans.clear()
|
||||
response = await ds.client.get("/no-such-database-at-all")
|
||||
assert response.status_code == 404
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["http.response.status_code"] == 404
|
||||
assert "error.type" not in server[0].attributes
|
||||
assert server[0].status.status_code is StatusCode.UNSET
|
||||
# Route enrichment must not be gated on a successful response.
|
||||
assert "http.route" in server[0].attributes
|
||||
assert server[0].name != "GET"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_an_unrouted_404_has_no_route_and_a_bare_method_name(ds, otel_spans):
|
||||
"""
|
||||
When no route matches there is nothing to set `http.route` to, so the span
|
||||
keeps the bare method name it was given at the edge - which is exactly the
|
||||
fallback semantic conventions specify for an unknown route.
|
||||
|
||||
`/a/b/c/d/e` is used rather than a plausible-looking missing name because
|
||||
Datasette's route table is greedy: `/no-such-database-at-all` matches the
|
||||
database pattern, and `/-/nope/deeper` matches the row pattern. Only a
|
||||
path deeper than any route matches nothing at all.
|
||||
"""
|
||||
otel_spans.clear()
|
||||
response = await ds.client.get("/a/b/c/d/e")
|
||||
assert response.status_code == 404
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].name == "GET"
|
||||
assert "http.route" not in server[0].attributes
|
||||
assert server[0].attributes["http.response.status_code"] == 404
|
||||
assert server[0].status.status_code is StatusCode.UNSET
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_only_the_first_http_response_start_is_recorded(otel_spans):
|
||||
"""
|
||||
The `send` wrapper keeps the first status it sees.
|
||||
|
||||
Nothing in Datasette sends two `http.response.start` messages, so this
|
||||
drives the middleware directly rather than pretending a request could
|
||||
reach it. Without the guard a misbehaving plugin's second start message
|
||||
would silently replace the status the client actually received.
|
||||
"""
|
||||
|
||||
async def two_starts(scope, receive, send):
|
||||
await send({"type": "http.response.start", "status": 200, "headers": []})
|
||||
await send({"type": "http.response.start", "status": 503, "headers": []})
|
||||
await send({"type": "http.response.body", "body": b""})
|
||||
|
||||
middleware = TelemetryMiddleware(two_starts)
|
||||
scope = {
|
||||
"type": "http",
|
||||
"method": "GET",
|
||||
"path": "/twice",
|
||||
"raw_path": b"/twice",
|
||||
"scheme": "http",
|
||||
"headers": [],
|
||||
}
|
||||
otel_spans.clear()
|
||||
await middleware(scope, None, lambda message: asyncio.sleep(0))
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["http.response.status_code"] == 200
|
||||
assert "error.type" not in server[0].attributes
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_lifespan_scope_passes_through_unspanned(otel_spans):
|
||||
"""
|
||||
`AsgiLifespan` sits *inside* this middleware, so the scope-type check has
|
||||
to come first or startup and shutdown events never reach it. A SERVER
|
||||
span for a lifespan scope is the symptom of that check being missing or
|
||||
late.
|
||||
"""
|
||||
instance = Datasette(memory=True)
|
||||
app = instance.app()
|
||||
events = iter([{"type": "lifespan.startup"}, {"type": "lifespan.shutdown"}])
|
||||
sent = []
|
||||
|
||||
async def receive():
|
||||
return next(events)
|
||||
|
||||
async def send(message):
|
||||
sent.append(message["type"])
|
||||
|
||||
otel_spans.clear()
|
||||
await app({"type": "lifespan"}, receive, send)
|
||||
assert sent == ["lifespan.startup.complete", "lifespan.shutdown.complete"]
|
||||
assert not _server_spans(otel_spans)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_http_route_is_the_compiled_pattern(ds, otel_spans):
|
||||
"""
|
||||
`http.route` is the route's compiled regex, not a prettified template.
|
||||
|
||||
Asserted against what Datasette's own router resolves rather than against
|
||||
a copied literal, so this pins the *relationship* - the attribute is the
|
||||
matched route - and does not break when a core pattern is edited.
|
||||
"""
|
||||
path = f"/{ds.db_name}/t"
|
||||
expected = _route_for(ds, path)
|
||||
otel_spans.clear()
|
||||
assert (await ds.client.get(path)).status_code == 200
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["http.route"] == expected
|
||||
assert server[0].name == f"GET {expected}"
|
||||
# The pattern really is the ugly one, and that is deliberate - if someone
|
||||
# adds a prettifier this is the assertion that should make them argue for
|
||||
# it rather than slip it in.
|
||||
assert "(?P<database>" in expected
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_the_route_lands_on_the_request_span_not_a_plugins_current_span(
|
||||
ds, otel_spans
|
||||
):
|
||||
"""
|
||||
The route is set on the span the middleware started, found through the
|
||||
ASGI scope - not on whatever span happens to be current when routing
|
||||
resolves.
|
||||
|
||||
Those are the same span only until a plugin `asgi_wrapper()` starts one of
|
||||
its own. A plugin wrapper runs *inside* this middleware, so an instrumented
|
||||
plugin makes its span current for the whole request: reading the current
|
||||
span in `route_path` renames that plugin's INTERNAL span to
|
||||
`GET <route>` and hangs `http.route` off it, while the actual request span
|
||||
keeps a bare method name and never gets the one attribute a trace UI
|
||||
groups requests by. Verified by reproducing it, not by reasoning about it.
|
||||
"""
|
||||
ds.pm.register(_MiddlewarePlugin(), name="httpspan-middleware")
|
||||
try:
|
||||
otel_spans.clear()
|
||||
path = f"/{ds.db_name}/t"
|
||||
expected = _route_for(ds, path)
|
||||
assert (await ds.client.get(path)).status_code == 200
|
||||
finally:
|
||||
ds.pm.unregister(name="httpspan-middleware")
|
||||
|
||||
spans = otel_spans.get_finished_spans()
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["http.route"] == expected
|
||||
assert server[0].name == f"GET {expected}"
|
||||
# And the plugin's span is untouched: same name, no route attribute.
|
||||
plugin_spans = [span for span in spans if span.name == PLUGIN_MIDDLEWARE_SPAN]
|
||||
assert len(plugin_spans) == 1
|
||||
assert "http.route" not in (plugin_spans[0].attributes or {})
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_request_span_attributes(ds, otel_spans):
|
||||
"The whole attribute set on one ordinary request."
|
||||
path = f"/{ds.db_name}/t"
|
||||
otel_spans.clear()
|
||||
assert (await ds.client.get(path)).status_code == 200
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
attributes = server[0].attributes
|
||||
assert attributes["http.request.method"] == "GET"
|
||||
assert attributes["url.path"] == path
|
||||
assert attributes["url.scheme"] == "http"
|
||||
assert attributes["http.response.status_code"] == 200
|
||||
assert attributes["http.route"] == _route_for(ds, path)
|
||||
assert server[0].status.status_code is StatusCode.UNSET
|
||||
# Never, on any span: an IP is borderline PII and the query string carries
|
||||
# user-supplied SQL.
|
||||
assert "client.address" not in attributes
|
||||
assert "url.query" not in attributes
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_db_query_spans_are_children_of_the_request_span(ds, otel_spans):
|
||||
"""
|
||||
The point of the whole PR.
|
||||
|
||||
Not just "same trace ID" - every `db.query` span must reach the request
|
||||
span by walking parents, and the request span must be the only root. A
|
||||
stray root would show up in a trace UI as its own single-span trace, which
|
||||
is the state this replaces.
|
||||
"""
|
||||
otel_spans.clear()
|
||||
assert (await ds.client.get(f"/{ds.db_name}/t?_facet=v")).status_code == 200
|
||||
spans = otel_spans.get_finished_spans()
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
server_span = server[0]
|
||||
assert server_span.parent is None
|
||||
|
||||
by_span_id = {span.context.span_id: span for span in spans}
|
||||
roots = [span for span in spans if span.parent is None]
|
||||
assert [span.name for span in roots] == [server_span.name], (
|
||||
"every span from a request should hang off the request span, but these "
|
||||
f"are roots: {sorted(span.name for span in roots)}"
|
||||
)
|
||||
|
||||
queries = [span for span in spans if span.name == "db.query"]
|
||||
assert queries, "a faceted table page should have issued queries"
|
||||
for query in queries:
|
||||
assert query.context.trace_id == server_span.context.trace_id
|
||||
# Walk up to the root, which must be the request span.
|
||||
current = query
|
||||
seen = 0
|
||||
while current.parent is not None:
|
||||
current = by_span_id[current.parent.span_id]
|
||||
seen += 1
|
||||
assert seen < 20, "parent chain did not terminate"
|
||||
assert current is server_span
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_500_sets_error_status_and_error_type(ds, otel_spans):
|
||||
"""
|
||||
A plain 500 - no exception escaping the app, because `route_path` converts
|
||||
it into a response itself. The status is the only signal the middleware
|
||||
gets, so `error.type` is the status as a string.
|
||||
"""
|
||||
ds.pm.register(_BoomPlugin(), name="httpspan-boom")
|
||||
try:
|
||||
otel_spans.clear()
|
||||
response = await ds.client.get("/-/http-span-boom")
|
||||
assert response.status_code == 500
|
||||
finally:
|
||||
ds.pm.unregister(name="httpspan-boom")
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
assert server[0].attributes["http.response.status_code"] == 500
|
||||
assert server[0].attributes["error.type"] == "500"
|
||||
assert server[0].status.status_code is StatusCode.ERROR
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_csv_stream_span_covers_the_body_send(ds_paging, otel_spans):
|
||||
"""
|
||||
The span must not end when the handler returns - it has to cover the
|
||||
response body.
|
||||
|
||||
`stream_csv` runs its generator inline inside `AsgiStream.asgi_send`, and
|
||||
that call happens inside the single `await self.app(...)` the middleware
|
||||
makes, so a plain `finally` is enough and no deferred-end machinery is
|
||||
needed. This is the assertion that holds that claim up: a `db.query` that
|
||||
starts during the body send must still finish before the request span
|
||||
does.
|
||||
|
||||
Only meaningful on an export that actually pages, hence `ds_paging` - on a
|
||||
single-page table every query is over before the body begins and this
|
||||
passes however early the span ends. The middle assertion below, that some
|
||||
query *started* after `http.response.start` went out, is what keeps the
|
||||
test honest about that; it is why the app is driven as raw ASGI rather
|
||||
than through `ds.client`, which cannot timestamp the response start.
|
||||
|
||||
`time.time_ns()` is the same clock the SDK stamps spans with, so the two
|
||||
are directly comparable.
|
||||
"""
|
||||
app = ds_paging.app()
|
||||
body = []
|
||||
response_started_at = None
|
||||
|
||||
async def receive():
|
||||
return {"type": "http.request", "body": b"", "more_body": False}
|
||||
|
||||
async def send(message):
|
||||
nonlocal response_started_at
|
||||
if message["type"] == "http.response.start":
|
||||
assert message["status"] == 200
|
||||
response_started_at = time.time_ns()
|
||||
else:
|
||||
body.append(message.get("body") or b"")
|
||||
|
||||
otel_spans.clear()
|
||||
await app(
|
||||
{
|
||||
"type": "http",
|
||||
"http_version": "1.1",
|
||||
"method": "GET",
|
||||
"path": f"/{ds_paging.db_name}/t.csv",
|
||||
"raw_path": f"/{ds_paging.db_name}/t.csv".encode("latin-1"),
|
||||
"query_string": b"_stream=1",
|
||||
"scheme": "http",
|
||||
"headers": [(b"host", b"localhost")],
|
||||
},
|
||||
receive,
|
||||
send,
|
||||
)
|
||||
# 40 rows plus a header - the export really did read past one page
|
||||
assert len(b"".join(body).decode("utf-8").strip().splitlines()) == 41
|
||||
assert response_started_at is not None
|
||||
|
||||
spans = otel_spans.get_finished_spans()
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
server_span = server[0]
|
||||
queries = [span for span in spans if span.name == "db.query"]
|
||||
assert len(queries) > 1
|
||||
during_body = [span for span in queries if span.start_time > response_started_at]
|
||||
assert during_body, (
|
||||
"no query ran after the response started, so this workload cannot "
|
||||
"distinguish a span that covers the body send from one that ends when "
|
||||
"the handler returns - the export is not paging"
|
||||
)
|
||||
last_query_end = max(span.end_time for span in queries)
|
||||
assert server_span.end_time > last_query_end, (
|
||||
"the request span ended before the last query of a streaming export - "
|
||||
"it is not covering the response body"
|
||||
)
|
||||
for query in queries:
|
||||
assert query.context.trace_id == server_span.context.trace_id
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_inbound_traceparent_becomes_the_parent(ds, otel_spans):
|
||||
"""
|
||||
W3C trace context is extracted with the global propagator, so a request
|
||||
from an already-traced caller continues that trace.
|
||||
|
||||
The sampled flag has to be set: the SDK's default sampler is
|
||||
parentbased_always_on, so a `-00` flag would drop the span and the test
|
||||
would fail for a reason that has nothing to do with propagation.
|
||||
"""
|
||||
trace_id = "4bf92f3577b34da6a3ce929d0e0e4736"
|
||||
parent_span_id = "00f067aa0ba902b7"
|
||||
otel_spans.clear()
|
||||
response = await ds.client.get(
|
||||
f"/{ds.db_name}/t",
|
||||
headers={"traceparent": f"00-{trace_id}-{parent_span_id}-01"},
|
||||
)
|
||||
assert response.status_code == 200
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
server_span = server[0]
|
||||
assert f"{server_span.context.trace_id:032x}" == trace_id
|
||||
assert server_span.parent is not None
|
||||
assert f"{server_span.parent.span_id:016x}" == parent_span_id
|
||||
assert server_span.parent.is_remote
|
||||
# And the database spans joined the caller's trace too, not a new one.
|
||||
queries = [
|
||||
span for span in otel_spans.get_finished_spans() if span.name == "db.query"
|
||||
]
|
||||
assert queries
|
||||
for query in queries:
|
||||
assert f"{query.context.trace_id:032x}" == trace_id
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_user_supplied_sql_in_the_query_string_is_never_recorded(ds, otel_spans):
|
||||
"""
|
||||
The `?sql=` case specifically, which is the one that matters: this is the
|
||||
request where the query string *is* user-supplied SQL, and it reaches a
|
||||
view that runs it. The marker is searched for across every attribute of
|
||||
every span in the trace, not just for a `url.query` key, so recording it
|
||||
under some other name fails too.
|
||||
|
||||
`db.query.text` legitimately contains the SQL - that is documented and
|
||||
deliberate - so the marker is checked against the request span's own
|
||||
attributes, and against `url.*` and `http.*` keys everywhere.
|
||||
"""
|
||||
marker = "secret_marker_5b1f"
|
||||
otel_spans.clear()
|
||||
# `/{db}?sql=` 302s to the query view, so go straight there - a redirect
|
||||
# would leave the SQL only on a span for a request that never ran it.
|
||||
response = await ds.client.get(f"/{ds.db_name}/-/query?sql=select+'{marker}'")
|
||||
assert response.status_code == 200
|
||||
spans = otel_spans.get_finished_spans()
|
||||
server = _server_spans(otel_spans)
|
||||
assert len(server) == 1
|
||||
leaked = [
|
||||
f"{span.name} -> {key}={value!r}"
|
||||
for span in spans
|
||||
for key, value in (span.attributes or {}).items()
|
||||
if (span is server[0] or str(key).startswith(("url.", "http.")))
|
||||
and (marker in str(value) or str(key) == "url.query")
|
||||
]
|
||||
assert not leaked, "the query string reached a span attribute: " + ", ".join(leaked)
|
||||
# The request really did carry the marker, so the search above had
|
||||
# something to find.
|
||||
assert marker in response.text
|
||||
|
||||
|
||||
def test_request_span_skips_a_valid_but_non_recording_span():
|
||||
"""
|
||||
`request_span()` is guarded on `is_recording()`, not on
|
||||
`get_span_context().is_valid`, and this is the case that separates them.
|
||||
|
||||
With no provider installed but an inbound `traceparent`, the API's
|
||||
NoOpTracer hands back a `NonRecordingSpan` carrying the *remote* span
|
||||
context - valid, sampled, and recording nothing. An `is_valid` guard would
|
||||
wave that through and the router would build the name string and call
|
||||
`set_attribute`/`update_name` on a span that discards both.
|
||||
|
||||
Tested at this level deliberately: through a real request the two guards
|
||||
are indistinguishable, because every call the router makes on a
|
||||
NonRecordingSpan is already a no-op. The only difference is the work done
|
||||
to get there, so the guard itself is what has to be asserted on.
|
||||
"""
|
||||
remote = SpanContext(
|
||||
trace_id=0x4BF92F3577B34DA6A3CE929D0E0E4736,
|
||||
span_id=0x00F067AA0BA902B7,
|
||||
is_remote=True,
|
||||
trace_flags=TraceFlags(TraceFlags.SAMPLED),
|
||||
)
|
||||
assert remote.is_valid
|
||||
non_recording = NonRecordingSpan(remote)
|
||||
assert non_recording.is_recording() is False
|
||||
assert request_span({REQUEST_SPAN_SCOPE_KEY: non_recording}) is None
|
||||
# Nothing current, nothing in the scope: the INVALID_SPAN fallback.
|
||||
assert request_span({}) is None
|
||||
# And the case it must not skip.
|
||||
with tracer.start_as_current_span("test.request_span.recording") as span:
|
||||
assert request_span({REQUEST_SPAN_SCOPE_KEY: span}) is span
|
||||
# Falling back to the current span is how an externally installed
|
||||
# SERVER span still gets enriched.
|
||||
assert request_span({}) is span
|
||||
|
||||
|
||||
NO_PROVIDER_PROGRAM = textwrap.dedent("""
|
||||
import asyncio, json, sys
|
||||
|
||||
from datasette.telemetry import TelemetryMiddleware
|
||||
|
||||
seen = {}
|
||||
|
||||
|
||||
async def inner(scope, receive, send):
|
||||
seen.setdefault("sends", []).append(send)
|
||||
seen.setdefault("scopes", []).append(scope)
|
||||
await send({"type": "http.response.start", "status": 200, "headers": []})
|
||||
await send({"type": "http.response.body", "body": b""})
|
||||
|
||||
|
||||
async def real_send(message):
|
||||
pass
|
||||
|
||||
|
||||
async def main():
|
||||
middleware = TelemetryMiddleware(inner)
|
||||
for headers in ([], [(b"traceparent", b"00-" + b"a" * 32 + b"-" + b"b" * 16 + b"-01")]):
|
||||
await middleware(
|
||||
{
|
||||
"type": "http",
|
||||
"method": "GET",
|
||||
"path": "/",
|
||||
"raw_path": b"/",
|
||||
"scheme": "http",
|
||||
"headers": headers,
|
||||
},
|
||||
None,
|
||||
real_send,
|
||||
)
|
||||
print(
|
||||
json.dumps(
|
||||
{
|
||||
"unwrapped": [send is real_send for send in seen["sends"]],
|
||||
"scope_keys": [
|
||||
"datasette.telemetry.request_span" in scope
|
||||
for scope in seen["scopes"]
|
||||
],
|
||||
"sdk_imported": any(
|
||||
name.startswith("opentelemetry.sdk") for name in sys.modules
|
||||
),
|
||||
}
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
asyncio.run(main())
|
||||
""")
|
||||
|
||||
|
||||
def test_no_provider_takes_the_fast_path():
|
||||
"""
|
||||
With no `TracerProvider` installed the middleware must hand the
|
||||
application the *original* `send`, not a wrapper - a default Datasette
|
||||
install should pay essentially nothing for instrumentation it is not
|
||||
using.
|
||||
|
||||
This has to run in a subprocess. The suite's `otel_provider` fixture is
|
||||
session-scoped and autouse, and `set_tracer_provider()` is effectively
|
||||
once-per-process, so in-process every span is recording and the fast path
|
||||
is unreachable.
|
||||
|
||||
The second case, with an inbound `traceparent`, is the one that pins the
|
||||
check itself. With no provider the API's NoOpTracer returns a
|
||||
NonRecordingSpan carrying the *remote* span context: its
|
||||
`get_span_context().is_valid` is True while `is_recording()` is False. A
|
||||
fast path guarded on `is_valid` would therefore silently stop working for
|
||||
exactly the requests that arrive from an already-traced caller - which on
|
||||
a real deployment behind an instrumented proxy is all of them.
|
||||
|
||||
conftest.py's pytest_collection_modifyitems() moves this test to the front
|
||||
of the run by name - if you rename it, rename it there too.
|
||||
"""
|
||||
result = subprocess.run(
|
||||
[sys.executable, "-c", NO_PROVIDER_PROGRAM],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=True,
|
||||
)
|
||||
report = json.loads(result.stdout)
|
||||
assert report["sdk_imported"] is False, "the SDK loaded in a fresh interpreter"
|
||||
assert report["unwrapped"] == [True, True], (
|
||||
"the middleware wrapped `send` with no provider installed; the second "
|
||||
"entry is the inbound-traceparent case, which fails if the fast path "
|
||||
"is guarded on is_valid instead of is_recording()"
|
||||
)
|
||||
# Same fast path, other observable: nothing is stashed in the scope either.
|
||||
assert report["scope_keys"] == [False, False]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_internal_client_requests_are_marked(ds, otel_spans):
|
||||
"""
|
||||
An in-process `datasette.client` request runs the full ASGI stack, so it
|
||||
emits its own SERVER span - `datasette.internal_client` marks those so
|
||||
kind-based dashboards can filter the double-count out. A request that
|
||||
arrives through the raw ASGI app (the shape of a real inbound request,
|
||||
without the DatasetteClient wrapper setting the ContextVar) must not
|
||||
carry the attribute.
|
||||
"""
|
||||
otel_spans.clear()
|
||||
assert (await ds.client.get("/")).status_code == 200
|
||||
server = _server_spans(otel_spans)
|
||||
assert server
|
||||
assert all(
|
||||
span.attributes.get("datasette.internal_client") is True for span in server
|
||||
)
|
||||
|
||||
import httpx
|
||||
|
||||
transport = httpx.ASGITransport(app=ds.app())
|
||||
async with httpx.AsyncClient(
|
||||
transport=transport, base_url="http://localhost"
|
||||
) as client:
|
||||
otel_spans.clear()
|
||||
assert (await client.get("/")).status_code == 200
|
||||
server = _server_spans(otel_spans)
|
||||
assert server
|
||||
assert all("datasette.internal_client" not in span.attributes for span in server)
|
||||
|
|
@ -3,13 +3,11 @@ Tests for the datasette.database.Database class
|
|||
"""
|
||||
|
||||
import asyncio
|
||||
import threading
|
||||
import uuid
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
import sqlite_utils
|
||||
from opentelemetry import context as otel_context_api
|
||||
|
||||
from datasette.app import Datasette
|
||||
from datasette.database import (
|
||||
|
|
@ -1225,110 +1223,3 @@ async def test_database_close_is_idempotent(tmpdir):
|
|||
# Second call should be a no-op, not raise
|
||||
db.close()
|
||||
ds._internal_database.close()
|
||||
|
||||
|
||||
_CONTEXT_LEAK_MARKER_KEY = "otel-context-leak-marker"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize("num_sql_threads", (0, 1))
|
||||
async def test_write_thread_context_is_detached_between_tasks(
|
||||
tmp_path, monkeypatch, num_sql_threads
|
||||
):
|
||||
"""
|
||||
The write thread attaches each task's otel Context and must detach it
|
||||
again before picking up the next task. The thread is persistent and
|
||||
shared, so a leaked token would grow that thread's context stack for the
|
||||
rest of the process - and a *wrong*-token detach only logs a warning
|
||||
rather than raising, so "does it throw" cannot catch either mistake.
|
||||
|
||||
Two things are asserted, because neither alone is sufficient:
|
||||
|
||||
1. Each task observes the context value that was current on the event
|
||||
loop when it was queued. This is what fails if the Context is not
|
||||
carried on WriteTask, or is never attached. It does *not* catch a
|
||||
missing detach: attach() replaces the current Context wholesale, so a
|
||||
leftover one from a previous task is simply overwritten.
|
||||
2. The write thread's attach depth is identical at the same point in
|
||||
every task. This is what fails if detach is missing - the stack grows
|
||||
by one per task - and it holds across a task that raises, because the
|
||||
detach lives in a `finally`.
|
||||
|
||||
An otel context value is used rather than a plain contextvars.ContextVar:
|
||||
a plain var set on the event loop never crosses into the write thread, so
|
||||
the probe would read None every time and the test could not fail.
|
||||
"""
|
||||
name = f"context_leak_test_{num_sql_threads}"
|
||||
db_path = tmp_path / f"{name}.db"
|
||||
sqlite3.connect(db_path).close()
|
||||
ds = Datasette([str(db_path)], settings={"num_sql_threads": num_sql_threads})
|
||||
db = ds.get_database(name)
|
||||
await db.execute_write("create table t (id integer primary key)")
|
||||
|
||||
write_thread_name = f"_execute_writes for database {name}"
|
||||
depth = {"value": 0}
|
||||
real_attach = otel_context_api.attach
|
||||
real_detach = otel_context_api.detach
|
||||
|
||||
def counting_attach(context):
|
||||
token = real_attach(context)
|
||||
if threading.current_thread().name == write_thread_name:
|
||||
depth["value"] += 1
|
||||
return token
|
||||
|
||||
def counting_detach(token):
|
||||
real_detach(token)
|
||||
if threading.current_thread().name == write_thread_name:
|
||||
depth["value"] -= 1
|
||||
|
||||
# Patched on the opentelemetry.context module itself, which is what both
|
||||
# database.py and opentelemetry.trace.use_span() look the functions up on.
|
||||
monkeypatch.setattr(otel_context_api, "attach", counting_attach)
|
||||
monkeypatch.setattr(otel_context_api, "detach", counting_detach)
|
||||
|
||||
seen_markers = []
|
||||
seen_depths = []
|
||||
|
||||
def probe(conn):
|
||||
seen_markers.append(otel_context_api.get_value(_CONTEXT_LEAK_MARKER_KEY))
|
||||
seen_depths.append(depth["value"])
|
||||
|
||||
def failing_probe(conn):
|
||||
probe(conn)
|
||||
# Exercises the write thread's exception path: the detach still has
|
||||
# to happen, which is why it lives in a `finally`.
|
||||
raise ValueError("deliberate failure inside a write task")
|
||||
|
||||
try:
|
||||
for i in range(5):
|
||||
ctx = otel_context_api.set_value(_CONTEXT_LEAK_MARKER_KEY, f"marker-{i}")
|
||||
token = real_attach(ctx)
|
||||
try:
|
||||
if i == 2:
|
||||
with pytest.raises(ValueError):
|
||||
await db.execute_write_fn(failing_probe)
|
||||
else:
|
||||
await db.execute_write_fn(probe)
|
||||
finally:
|
||||
real_detach(token)
|
||||
|
||||
# Sanity check: no marker is active in *this* (event loop) context
|
||||
# right now, so the final probe is a fair test of the write thread's
|
||||
# own state rather than something this test forgot to clean up.
|
||||
assert otel_context_api.get_value(_CONTEXT_LEAK_MARKER_KEY) is None
|
||||
await db.execute_write_fn(probe)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
assert seen_markers == [
|
||||
"marker-0",
|
||||
"marker-1",
|
||||
"marker-2",
|
||||
"marker-3",
|
||||
"marker-4",
|
||||
None,
|
||||
]
|
||||
assert len(set(seen_depths)) == 1, (
|
||||
f"write thread context stack grew across tasks: {seen_depths} - "
|
||||
"a token was attached without being detached"
|
||||
)
|
||||
|
|
|
|||
File diff suppressed because it is too large
Load diff
|
|
@ -1,532 +0,0 @@
|
|||
"""
|
||||
Tests for the OpenTelemetry metrics Datasette core emits.
|
||||
|
||||
Two layers are tested separately and deliberately:
|
||||
|
||||
- The gauge callbacks are plain generator functions, so they are called
|
||||
directly for exact-value assertions. Going through the SDK for those would
|
||||
be unreliable: the pool gauges carry no attribute identifying which
|
||||
Datasette produced them, and a pytest session has many live instances, so
|
||||
the SDK's last-value aggregation would report whichever one happened to be
|
||||
observed last.
|
||||
|
||||
- The SDK pipeline (instrument -> reader -> data points) is tested through
|
||||
the `otel_metrics` fixture, using metrics that carry `db.namespace` - a
|
||||
uniquely named in-memory database is enough to isolate those from every
|
||||
other instance alive in the session.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import threading
|
||||
import weakref
|
||||
|
||||
import pytest
|
||||
|
||||
from datasette import telemetry
|
||||
from datasette.app import Datasette
|
||||
from datasette.database import Database
|
||||
from datasette.utils.sqlite import sqlite3
|
||||
|
||||
pytestmark = pytest.mark.filterwarnings("ignore::ResourceWarning")
|
||||
|
||||
|
||||
def observations(callback, datasette=None):
|
||||
"""
|
||||
Run a gauge callback, optionally keeping only observations produced by one
|
||||
Datasette's databases. Returns a list of (attributes dict, value).
|
||||
"""
|
||||
results = []
|
||||
names = None
|
||||
if datasette is not None:
|
||||
names = {db.name for db in telemetry._databases_of(datasette)}
|
||||
for observation in callback():
|
||||
attributes = dict(observation.attributes or {})
|
||||
namespace = attributes.get("db.namespace")
|
||||
if names is not None and namespace is not None and namespace not in names:
|
||||
continue
|
||||
results.append((attributes, observation.value))
|
||||
return results
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def metrics_ds():
|
||||
"A Datasette with a distinctive thread count and a uniquely named database."
|
||||
ds = Datasette(
|
||||
memory=True,
|
||||
settings={"num_sql_threads": 7},
|
||||
)
|
||||
ds.add_memory_database("metrics_test_db")
|
||||
try:
|
||||
yield ds
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_sql_thread_limit_gauge_reports_num_sql_threads(metrics_ds):
|
||||
values = [value for _, value in observations(telemetry.observe_sql_thread_limit)]
|
||||
# Other instances are alive in this session, so assert membership rather
|
||||
# than uniqueness - 7 is distinctive enough to only come from metrics_ds.
|
||||
assert 7 in values
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_no_thread_gauges_in_non_threaded_mode():
|
||||
"""
|
||||
num_sql_threads=0 means there is no pool at all, so the pool gauges must
|
||||
skip the instance rather than report a bogus limit of 0.
|
||||
|
||||
The pool gauges carry no attributes, so the live-instance registry is
|
||||
narrowed to just this instance for the assertion - counting global
|
||||
observations instead would let an unrelated instance being garbage
|
||||
collected mid-test shift the baseline.
|
||||
"""
|
||||
ds = Datasette(memory=True, settings={"num_sql_threads": 0})
|
||||
try:
|
||||
assert ds.executor is None
|
||||
original = telemetry._live_datasettes
|
||||
telemetry._live_datasettes = weakref.WeakSet([ds])
|
||||
try:
|
||||
assert list(telemetry.observe_sql_thread_limit()) == []
|
||||
assert list(telemetry.observe_sql_thread_queue_depth()) == []
|
||||
finally:
|
||||
telemetry._live_datasettes = original
|
||||
# Per-database gauges are unaffected - they do not depend on the pool.
|
||||
assert observations(telemetry.observe_pending_queries, ds)
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_thread_queue_depth_gauge_reports_saturation():
|
||||
"""
|
||||
The headline alerting metric must actually read above zero when reads
|
||||
queue behind num_sql_threads. This also pins the private
|
||||
`ThreadPoolExecutor._work_queue` attribute the callback depends on: if a
|
||||
stdlib rename ever removes it, this fails instead of the metric silently
|
||||
vanishing (the callback tolerates its absence at collection time).
|
||||
"""
|
||||
ds = Datasette(memory=True, settings={"num_sql_threads": 1})
|
||||
db = ds.add_memory_database("metrics_saturation_db")
|
||||
entered = threading.Event()
|
||||
release = threading.Event()
|
||||
|
||||
def blocker(conn):
|
||||
entered.set()
|
||||
assert release.wait(timeout=10)
|
||||
return 1
|
||||
|
||||
try:
|
||||
first = asyncio.ensure_future(db.execute_fn(blocker))
|
||||
# Wait until the blocker owns the pool's only thread.
|
||||
await asyncio.get_running_loop().run_in_executor(None, entered.wait, 10)
|
||||
second = asyncio.ensure_future(db.execute_fn(lambda conn: 2))
|
||||
# The second submission lands in the executor's queue on the next
|
||||
# event-loop turn; poll briefly rather than assume the timing.
|
||||
depths = []
|
||||
for _ in range(500):
|
||||
depths = [
|
||||
value
|
||||
for _, value in observations(telemetry.observe_sql_thread_queue_depth)
|
||||
]
|
||||
if any(value >= 1 for value in depths):
|
||||
break
|
||||
await asyncio.sleep(0.01)
|
||||
assert any(value >= 1 for value in depths), depths
|
||||
release.set()
|
||||
assert await first == 1
|
||||
assert await second == 2
|
||||
finally:
|
||||
release.set()
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_pending_queries_gauge_tracks_in_flight_queries(metrics_ds):
|
||||
db = metrics_ds.get_database("metrics_test_db")
|
||||
attributes = {"db.namespace": "metrics_test_db"}
|
||||
|
||||
def value():
|
||||
points = [
|
||||
v
|
||||
for a, v in observations(telemetry.observe_pending_queries, metrics_ds)
|
||||
if a == attributes
|
||||
]
|
||||
assert len(points) == 1
|
||||
return points[0]
|
||||
|
||||
assert value() == 0
|
||||
|
||||
# sqlite3.sleep is not a thing, so block the worker thread on an event we
|
||||
# control from the event loop and sample the gauge while it is held.
|
||||
release = asyncio.Event()
|
||||
loop = asyncio.get_running_loop()
|
||||
entered = asyncio.Event()
|
||||
|
||||
def blocking_fn(conn):
|
||||
loop.call_soon_threadsafe(entered.set)
|
||||
asyncio.run_coroutine_threadsafe(release.wait(), loop).result()
|
||||
return "done"
|
||||
|
||||
task = asyncio.ensure_future(db.execute_fn(blocking_fn))
|
||||
await entered.wait()
|
||||
assert value() == 1, "a query occupying a pool thread must be counted as pending"
|
||||
release.set()
|
||||
assert await task == "done"
|
||||
assert value() == 0, "the count must drop once the query completes"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_write_queue_depth_gauge(metrics_ds):
|
||||
db = metrics_ds.get_database("metrics_test_db")
|
||||
attributes = {"db.namespace": "metrics_test_db"}
|
||||
|
||||
def depths():
|
||||
return [
|
||||
v
|
||||
for a, v in observations(telemetry.observe_write_queue_depth, metrics_ds)
|
||||
if a == attributes
|
||||
]
|
||||
|
||||
# No write has ever been queued, so there is no queue and no observation -
|
||||
# rather than a fabricated zero for a queue that does not exist.
|
||||
assert depths() == []
|
||||
|
||||
await db.execute_write("create table t (id integer primary key)")
|
||||
assert depths() == [0], "an idle write queue reports zero, not nothing"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_open_connections_gauge(metrics_ds, tmp_path):
|
||||
path = str(tmp_path / "conns.db")
|
||||
sqlite3.connect(path).execute("create table t (id integer primary key)")
|
||||
db = metrics_ds.add_database(Database(metrics_ds, path=path), name="conns_db")
|
||||
attributes = {"db.namespace": "conns_db"}
|
||||
|
||||
def open_connections():
|
||||
points = [
|
||||
v
|
||||
for a, v in observations(telemetry.observe_open_connections, metrics_ds)
|
||||
if a == attributes
|
||||
]
|
||||
assert len(points) == 1
|
||||
return points[0]
|
||||
|
||||
assert open_connections() == 0
|
||||
await db.execute("select 1")
|
||||
assert open_connections() >= 1, "executing a query opens a tracked file connection"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_operation_duration_histogram_read(otel_metrics):
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("duration_read_db")
|
||||
try:
|
||||
db = ds.get_database("duration_read_db")
|
||||
await db.execute("select 1")
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "duration_read_db", "datasette.operation": "read"},
|
||||
)
|
||||
assert point.count == 1
|
||||
assert point.sum > 0
|
||||
assert dict(point.attributes)["db.system"] == "sqlite"
|
||||
assert "error.type" not in dict(point.attributes)
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_operation_duration_histogram_write(otel_metrics):
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("duration_write_db")
|
||||
try:
|
||||
db = ds.get_database("duration_write_db")
|
||||
await db.execute_write("create table t (id integer primary key)")
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "duration_write_db", "datasette.operation": "write"},
|
||||
)
|
||||
assert point.count == 1
|
||||
assert point.sum > 0
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_operation_duration_records_error_type(otel_metrics):
|
||||
"A failed query is still timed, and is separable from a successful one."
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("duration_error_db")
|
||||
try:
|
||||
db = ds.get_database("duration_error_db")
|
||||
with pytest.raises(sqlite3.OperationalError):
|
||||
await db.execute("select * from nope")
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "duration_error_db", "datasette.operation": "read"},
|
||||
)
|
||||
assert point.count == 1
|
||||
assert dict(point.attributes)["error.type"] == "OperationalError"
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_operation_duration_records_write_error_type(otel_metrics):
|
||||
"""
|
||||
Same as the read-path error test, but the write wrappers time a different
|
||||
code path - `execute_write_fn`, the write thread and its reply future -
|
||||
so error propagation through them is pinned separately.
|
||||
"""
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("duration_write_error_db")
|
||||
try:
|
||||
db = ds.get_database("duration_write_error_db")
|
||||
with pytest.raises(sqlite3.OperationalError):
|
||||
await db.execute_write("insert into nope values (1)")
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "duration_write_error_db", "datasette.operation": "write"},
|
||||
)
|
||||
assert point.count == 1
|
||||
assert dict(point.attributes)["error.type"] == "OperationalError"
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_write_queue_wait_histogram(otel_metrics):
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("queue_wait_db")
|
||||
try:
|
||||
db = ds.get_database("queue_wait_db")
|
||||
await db.execute_write("create table t (id integer primary key)")
|
||||
await db.execute_write("insert into t (id) values (1)")
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"datasette.write.queue_wait", {"db.namespace": "queue_wait_db"}
|
||||
)
|
||||
assert point.count == 2, "one measurement per write dequeued"
|
||||
assert point.sum >= 0
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_interrupted_queries_counter(otel_metrics):
|
||||
"The count of time-limit kills, which sampled traces cannot provide."
|
||||
ds = Datasette(memory=True, settings={"sql_time_limit_ms": 1})
|
||||
ds.add_memory_database("interrupted_db")
|
||||
try:
|
||||
db = ds.get_database("interrupted_db")
|
||||
from datasette.database import QueryInterrupted
|
||||
|
||||
with pytest.raises(QueryInterrupted):
|
||||
await db.execute("""
|
||||
with recursive counter(x) as (
|
||||
select 0 union all select x + 1 from counter
|
||||
)
|
||||
select * from counter
|
||||
""")
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"datasette.sql.queries.interrupted", {"db.namespace": "interrupted_db"}
|
||||
)
|
||||
assert point.value == 1
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_metrics_are_reported_through_the_sdk_for_gauges(otel_metrics):
|
||||
"End-to-end: a gauge callback reaches the reader as a data point."
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("gauge_pipeline_db")
|
||||
try:
|
||||
await ds.get_database("gauge_pipeline_db").execute("select 1")
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"datasette.sql.queries.pending", {"db.namespace": "gauge_pipeline_db"}
|
||||
)
|
||||
assert point.value == 0
|
||||
assert otel_metrics.points("datasette.sql.threads.limit")
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
def test_closed_datasette_stops_being_observed():
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("closed_db")
|
||||
assert observations(telemetry.observe_pending_queries, ds)
|
||||
ds.close()
|
||||
names = [
|
||||
attributes.get("db.namespace")
|
||||
for attributes, _ in observations(telemetry.observe_pending_queries)
|
||||
]
|
||||
assert "closed_db" not in names
|
||||
|
||||
|
||||
def test_registry_holds_instances_weakly():
|
||||
"""
|
||||
Registering an instance must never be the thing that keeps it alive.
|
||||
|
||||
A stand-in object is used rather than a real Datasette because a Datasette
|
||||
with a temp-disk internal database is pinned for the life of the process
|
||||
by `Database.__init__`'s `atexit.register(self._cleanup_temp_file)`, which
|
||||
holds the Database, which holds the Datasette. That is pre-existing and
|
||||
unrelated to telemetry; what is tested here is that this registry adds no
|
||||
reference of its own.
|
||||
"""
|
||||
import gc
|
||||
import weakref
|
||||
|
||||
class FakeDatasette:
|
||||
pass
|
||||
|
||||
fake = FakeDatasette()
|
||||
telemetry.register_datasette(fake)
|
||||
assert fake in telemetry._live_instances()
|
||||
ref = weakref.ref(fake)
|
||||
del fake
|
||||
gc.collect()
|
||||
assert ref() is None
|
||||
assert not any(isinstance(ds, FakeDatasette) for ds in telemetry._live_instances())
|
||||
|
||||
|
||||
HISTOGRAM_PROBES = [
|
||||
# (instrument attribute on telemetry, metric name, isolating attributes)
|
||||
(
|
||||
"sql_operation_duration",
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "bucket_probe_operation"},
|
||||
),
|
||||
(
|
||||
"write_queue_wait",
|
||||
"datasette.write.queue_wait",
|
||||
{"db.namespace": "bucket_probe_queue_wait"},
|
||||
),
|
||||
]
|
||||
|
||||
# One value inside each of six distinct registry buckets. Under OpenTelemetry's
|
||||
# default boundaries - [0, 5, 10, 25, ...], meant for milliseconds - the first
|
||||
# five of these all land in (0, 5] and only 7.0 lands elsewhere, so the
|
||||
# "occupies six buckets" assertion below fails if the advisory is ever dropped.
|
||||
SPREAD = [0.00005, 0.0003, 0.002, 0.03, 0.8, 7.0]
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"instrument_name,metric_name,attributes",
|
||||
HISTOGRAM_PROBES,
|
||||
ids=[metric for _, metric, _ in HISTOGRAM_PROBES],
|
||||
)
|
||||
def test_histograms_spread_values_across_buckets(
|
||||
otel_metrics, instrument_name, metric_name, attributes
|
||||
):
|
||||
"""
|
||||
The registry's boundaries reach the SDK, and a realistic spread of
|
||||
seconds-scale durations occupies more than one bucket.
|
||||
|
||||
Recording onto the instrument directly rather than driving a workload is
|
||||
deliberate: real durations here are all tens of microseconds and would
|
||||
share a bucket no matter what the boundaries were, which is exactly the
|
||||
situation this test exists to detect.
|
||||
|
||||
`explicit_bounds` is compared against the registry rather than against the
|
||||
instrument's own configuration - the instrument is built *from* the
|
||||
registry, so that comparison would be a value against itself. What is
|
||||
checked here is that the advisory survived the trip through the SDK.
|
||||
"""
|
||||
from datasette.telemetry_registry import METRICS
|
||||
|
||||
metric = next(m for m in METRICS if m == metric_name)
|
||||
instrument = getattr(telemetry, instrument_name)
|
||||
for value in SPREAD:
|
||||
instrument.record(value, attributes)
|
||||
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(metric_name, attributes)
|
||||
|
||||
assert (
|
||||
tuple(point.explicit_bounds) == metric.buckets
|
||||
), "the registry's boundaries did not reach the SDK"
|
||||
assert point.count == len(SPREAD)
|
||||
occupied = [count for count in point.bucket_counts if count]
|
||||
assert len(occupied) == len(SPREAD), (
|
||||
f"expected each of {SPREAD} in its own bucket, got bucket counts "
|
||||
f"{list(point.bucket_counts)} for bounds {list(point.explicit_bounds)}"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_operation_duration_histogram_records_execute_fn(otel_metrics):
|
||||
"Callback-style reads land in the same histogram as SQL-string reads."
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("duration_fn_db")
|
||||
try:
|
||||
db = ds.get_database("duration_fn_db")
|
||||
|
||||
def read_one(conn):
|
||||
return conn.execute("select 1").fetchone()[0]
|
||||
|
||||
assert await db.execute_fn(read_one) == 1
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "duration_fn_db", "datasette.operation": "read"},
|
||||
)
|
||||
assert point.count == 1
|
||||
assert point.sum > 0
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_operation_duration_histogram_records_execute_write_fn(otel_metrics):
|
||||
"Callback-style writes - the JSON write API's whole diet - are counted too."
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("duration_write_fn_db")
|
||||
try:
|
||||
db = ds.get_database("duration_write_fn_db")
|
||||
|
||||
def create_table(conn):
|
||||
conn.execute("create table t (id integer primary key)")
|
||||
|
||||
await db.execute_write_fn(create_table)
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "duration_write_fn_db", "datasette.operation": "write"},
|
||||
)
|
||||
assert point.count == 1
|
||||
assert point.sum > 0
|
||||
finally:
|
||||
ds.close()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_operation_duration_records_callback_error_type(otel_metrics):
|
||||
"A callback that raises is still timed, with error.type from the exception."
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database("duration_fn_error_db")
|
||||
try:
|
||||
db = ds.get_database("duration_fn_error_db")
|
||||
|
||||
def boom(conn):
|
||||
raise ValueError("callback failed")
|
||||
|
||||
with pytest.raises(ValueError):
|
||||
await db.execute_fn(boom)
|
||||
otel_metrics.collect()
|
||||
point = otel_metrics.point(
|
||||
"db.client.operation.duration",
|
||||
{"db.namespace": "duration_fn_error_db", "datasette.operation": "read"},
|
||||
)
|
||||
assert point.count == 1
|
||||
assert dict(point.attributes)["error.type"] == "ValueError"
|
||||
finally:
|
||||
ds.close()
|
||||
|
|
@ -1,597 +0,0 @@
|
|||
"""
|
||||
Two-way conformance between `datasette/telemetry_registry.py` and what
|
||||
Datasette actually emits.
|
||||
|
||||
This is the test that makes the generated documentation trustworthy. cog
|
||||
guarantees the docs match the registry; this guarantees the registry matches
|
||||
the code. Without it, both could agree with each other and be wrong.
|
||||
|
||||
It checks both directions, and the second one is the one nothing else catches:
|
||||
|
||||
- **emitted but not registered** - instrumentation was added without
|
||||
documenting it, so the reference page silently omits it.
|
||||
- **registered but never emitted** - the reference page describes a span or
|
||||
attribute that no longer exists, which is worse than omitting it, because a
|
||||
reader will build a dashboard on it.
|
||||
|
||||
Both of those directions compare the code against the registry. Neither can
|
||||
catch a *rename*, because the call sites now take their names from the
|
||||
registry - move `DB_NAMESPACE` to `"db.namespace2"` and code and registry
|
||||
still agree with each other, while every existing dashboard breaks. So the
|
||||
literal names live here too, spelled out, and are asserted against both the
|
||||
registry and the wire. That is the one comparison in this file that is not
|
||||
made against a value derived from the registry itself.
|
||||
"""
|
||||
|
||||
import itertools
|
||||
|
||||
import pytest
|
||||
import pytest_asyncio
|
||||
|
||||
pytest.importorskip("opentelemetry.sdk")
|
||||
|
||||
from opentelemetry.trace import SpanKind
|
||||
|
||||
from datasette import hookimpl
|
||||
from datasette import telemetry_registry as reg
|
||||
from datasette.telemetry_testing import assert_metrics_conform, assert_metrics_covered
|
||||
from datasette.app import Datasette
|
||||
from datasette.database import QueryInterrupted
|
||||
from datasette.utils.sqlite import sqlite3
|
||||
|
||||
# The names as they appear on the wire, written out rather than read from the
|
||||
# registry. If a change to the registry makes one of these fail, that change
|
||||
# is renaming something a user's dashboards and saved queries depend on -
|
||||
# which is a decision to take deliberately, here, not a line to re-derive.
|
||||
EXPECTED_ATTRIBUTES = {
|
||||
"db.query": {
|
||||
"db.system",
|
||||
"db.namespace",
|
||||
"db.query.text",
|
||||
"datasette.callback",
|
||||
"db.operation.name",
|
||||
"db.collection.name",
|
||||
"datasette.param_count",
|
||||
"datasette.param_sets",
|
||||
"datasette.time_limit_ms",
|
||||
"datasette.rows_returned",
|
||||
"datasette.truncated",
|
||||
"datasette.interrupted",
|
||||
"datasette.sql_error_suppressed",
|
||||
"datasette.executescript",
|
||||
"datasette.executemany",
|
||||
},
|
||||
"db.query.execute": set(),
|
||||
"db.write.queue_wait": set(),
|
||||
"db.write.execute": {
|
||||
"datasette.isolated_connection",
|
||||
"datasette.transaction",
|
||||
},
|
||||
"datasette.startup": set(),
|
||||
}
|
||||
EXPECTED_SPANS = set(EXPECTED_ATTRIBUTES)
|
||||
|
||||
# The HTTP request span is handled separately because its name is composed at
|
||||
# runtime - the request method, then the route it matched - so there is no
|
||||
# fixed string to pin it to. What can still be pinned, and is what a dashboard
|
||||
# depends on, is the shape of that name and the attribute keys.
|
||||
#
|
||||
# The route half is deliberately not spelled out as a literal: it is a core
|
||||
# route regex, and pinning those here would make an unrelated routing change
|
||||
# fail the telemetry conformance test. What is pinned instead is that the name
|
||||
# is exactly the method, a space, and the span's own `http.route` value - the
|
||||
# `{method} {route}` shape semantic conventions specify. The workload below
|
||||
# only issues GETs, so a change that stopped clamping the method, or that
|
||||
# started naming the span after the path, fails here.
|
||||
EXPECTED_HTTP_SPAN_NAME = "{http.request.method} {http.route}"
|
||||
EXPECTED_HTTP_METHOD_NAMES = {"GET"}
|
||||
EXPECTED_HTTP_ATTRIBUTES = {
|
||||
"http.request.method",
|
||||
"http.route",
|
||||
"url.path",
|
||||
"url.scheme",
|
||||
"server.address",
|
||||
"user_agent.original",
|
||||
"http.response.status_code",
|
||||
"error.type",
|
||||
"datasette.internal_client",
|
||||
}
|
||||
|
||||
# The registry's own name for the request span is that template, not anything
|
||||
# that appears on the wire.
|
||||
EXPECTED_REGISTRY_ATTRIBUTES = dict(
|
||||
EXPECTED_ATTRIBUTES, **{EXPECTED_HTTP_SPAN_NAME: EXPECTED_HTTP_ATTRIBUTES}
|
||||
)
|
||||
EXPECTED_REGISTRY_NAMES = set(EXPECTED_REGISTRY_ATTRIBUTES)
|
||||
|
||||
# Named in-memory databases are shared-cache, so two Datasette instances using
|
||||
# the same name share one SQLite database - and the second `create table`
|
||||
# fails. Every workload below therefore gets its own name.
|
||||
_names = itertools.count()
|
||||
|
||||
|
||||
def _unique(prefix):
|
||||
return f"{prefix}{next(_names)}"
|
||||
|
||||
|
||||
class _BoomPlugin:
|
||||
"""
|
||||
A route that raises.
|
||||
|
||||
`error.type` on the request span is only ever set by a 5xx, and nothing
|
||||
in Datasette returns one on a healthy instance - `route_path` converts
|
||||
exceptions into a 500 itself, so the workload has to supply the
|
||||
exception.
|
||||
"""
|
||||
|
||||
__name__ = "TelemetryRegistryBoomPlugin"
|
||||
|
||||
@hookimpl
|
||||
def register_routes(self):
|
||||
return [(r"^/-/telemetry-registry-boom$", lambda: 1 / 0)]
|
||||
|
||||
|
||||
async def exercise():
|
||||
"""
|
||||
Drive enough of Datasette to emit every span and attribute the registry
|
||||
claims exists.
|
||||
|
||||
Each call is here because it is the only thing that produces some span or
|
||||
attribute - see the comments. If you add instrumentation on a path this
|
||||
does not reach, add the path rather than loosening the assertions.
|
||||
|
||||
Returns the instance so the caller can close it; startup happens inside
|
||||
so that the `datasette.startup` span lands in the collected set.
|
||||
"""
|
||||
name = _unique("registry")
|
||||
ds = Datasette(memory=True)
|
||||
ds.add_memory_database(name)
|
||||
# datasette.startup - and the internal catalog work nested under it
|
||||
await ds.invoke_startup()
|
||||
db = ds.get_database(name)
|
||||
|
||||
# Writes: db.write.queue_wait, db.write.execute, db.query
|
||||
await db.execute_write("create table t (id integer primary key, v text)")
|
||||
# datasette.executemany, datasette.param_sets
|
||||
await db.execute_write_many(
|
||||
"insert into t (id, v) values (?, ?)", [[i, f"v{i}"] for i in range(30)]
|
||||
)
|
||||
# datasette.executescript
|
||||
await db.execute_write_script("create table t2 (id integer); drop table t2;")
|
||||
# datasette.transaction=False - VACUUM cannot run inside a transaction
|
||||
await db.execute_write("vacuum", transaction=False)
|
||||
# datasette.isolated_connection=True
|
||||
await db.execute_isolated_fn(lambda conn: conn.execute("select 1").fetchone())
|
||||
|
||||
# datasette.callback, with named functions so the conformance run sees the
|
||||
# attribute's documented value shape (a qualname, not just "<lambda>")
|
||||
def registry_read_callback(conn):
|
||||
return conn.execute("select count(*) from t").fetchone()
|
||||
|
||||
def registry_write_callback(conn):
|
||||
conn.execute("insert into t (id, v) values (100, 'callback')")
|
||||
|
||||
await db.execute_fn(registry_read_callback)
|
||||
await db.execute_write_fn(registry_write_callback)
|
||||
|
||||
# Reads: db.query.execute, datasette.rows_returned, datasette.truncated,
|
||||
# datasette.param_count, datasette.time_limit_ms
|
||||
await db.execute("select * from t where id > :n", {"n": 5})
|
||||
await db.execute("select * from t", truncate=True)
|
||||
|
||||
# datasette.sql_error_suppressed - the caller is probing and treats
|
||||
# failure as an expected answer
|
||||
with pytest.raises(sqlite3.OperationalError):
|
||||
await db.execute("select nope from t", log_sql_errors=False)
|
||||
|
||||
# datasette.interrupted - only ever set when a query exceeds its time
|
||||
# limit, so the workload has to force one rather than exempt it. An
|
||||
# unbounded recursive CTE cannot finish, so 1ms is always exceeded.
|
||||
with pytest.raises(QueryInterrupted):
|
||||
await db.execute(
|
||||
"with recursive c(x) as (select 0 union all select x+1 from c) "
|
||||
"select * from c",
|
||||
custom_time_limit=1,
|
||||
)
|
||||
|
||||
# db.collection.name - set only by views that already know their table.
|
||||
# These requests are also what produces the HTTP request span and its
|
||||
# http.request.method / url.path / url.scheme / server.address /
|
||||
# user_agent.original / http.response.status_code attributes.
|
||||
assert (await ds.client.get(f"/{name}/t?_facet=v")).status_code == 200
|
||||
assert (await ds.client.get(f"/{name}/t/1.json")).status_code == 200
|
||||
|
||||
# error.type on the request span, which only a 5xx sets
|
||||
ds.pm.register(_BoomPlugin(), name="telemetry-registry-boom")
|
||||
try:
|
||||
response = await ds.client.get("/-/telemetry-registry-boom")
|
||||
assert response.status_code == 500
|
||||
finally:
|
||||
ds.pm.unregister(name="telemetry-registry-boom")
|
||||
return ds
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def emitted(otel_spans):
|
||||
"""
|
||||
Every (span name, span kind, attributes) triple a broad workload emits.
|
||||
|
||||
The kind is carried because the request span's name is composed at
|
||||
runtime, so `span_for()` resolves it by kind instead. The attributes are
|
||||
carried as a mapping rather than a set of keys because the request span's
|
||||
name has to be checked against its own `http.route` value.
|
||||
"""
|
||||
# otel_spans has already cleared the exporter, and nothing is cleared
|
||||
# after this point: the workload's own startup emits datasette.startup.
|
||||
ds = await exercise()
|
||||
spans = otel_spans.get_finished_spans()
|
||||
assert spans, "no spans captured - the fixture is not exercising anything"
|
||||
# str() because span.name is the registry's SpanName instance, and a set
|
||||
# of those would compare equal to literals but read confusingly in a
|
||||
# failure message.
|
||||
collected = tuple(
|
||||
(
|
||||
str(span.name),
|
||||
span.kind,
|
||||
{str(key): value for key, value in (span.attributes or {}).items()},
|
||||
)
|
||||
for span in spans
|
||||
)
|
||||
ds.close()
|
||||
return collected
|
||||
|
||||
|
||||
def _partition(emitted):
|
||||
"The statically named spans, and the dynamically named request spans."
|
||||
static = [record for record in emitted if record[1] is not SpanKind.SERVER]
|
||||
server = [record for record in emitted if record[1] is SpanKind.SERVER]
|
||||
return static, server
|
||||
|
||||
|
||||
def _keys_by_span(records):
|
||||
by_span = {}
|
||||
for name, _kind, attributes in records:
|
||||
by_span.setdefault(name, set()).update(attributes)
|
||||
return by_span
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_workload_emits_exactly_the_expected_names(emitted):
|
||||
"""
|
||||
The wire format, pinned to literals.
|
||||
|
||||
Not derived from the registry, so this is what catches a rename that the
|
||||
registry and the call sites make together.
|
||||
"""
|
||||
static, server = _partition(emitted)
|
||||
by_span = _keys_by_span(static)
|
||||
assert set(by_span) == EXPECTED_SPANS
|
||||
assert by_span == EXPECTED_ATTRIBUTES
|
||||
|
||||
assert server, "the workload made HTTP requests but no SERVER span was emitted"
|
||||
union = set()
|
||||
methods = set()
|
||||
for name, _kind, attributes in server:
|
||||
union |= set(attributes)
|
||||
route = attributes.get("http.route")
|
||||
# Every request in the workload matches a route, so every one of these
|
||||
# names must be `{method} {route}`. A 404 would be a bare method - the
|
||||
# http_route tests cover that case with a real request.
|
||||
assert route, f"the request span {name!r} carries no http.route"
|
||||
method, _, name_route = name.partition(" ")
|
||||
assert name_route == route, (
|
||||
f"the request span is named {name!r}, which is not the "
|
||||
f"`{{method}} {{route}}` of {method!r} and {route!r}"
|
||||
)
|
||||
methods.add(method)
|
||||
assert methods == EXPECTED_HTTP_METHOD_NAMES
|
||||
assert union == EXPECTED_HTTP_ATTRIBUTES
|
||||
|
||||
|
||||
def test_registry_matches_the_expected_names():
|
||||
"The other half of the rename check: the registry against the same literals."
|
||||
assert {str(span) for span in reg.SPANS} == EXPECTED_REGISTRY_NAMES
|
||||
for span in reg.SPANS:
|
||||
assert {
|
||||
str(attribute) for attribute in span.attributes
|
||||
} == EXPECTED_REGISTRY_ATTRIBUTES[str(span)], f"{span} attributes have drifted"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_every_emitted_span_is_registered(emitted):
|
||||
"A span added without a registry entry would be missing from the docs."
|
||||
unregistered = sorted(
|
||||
{name for name, kind, _ in emitted if reg.span_for(name, kind) is None}
|
||||
)
|
||||
assert (
|
||||
not unregistered
|
||||
), f"these spans are emitted but not in telemetry_registry.SPANS: {unregistered}"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_every_emitted_attribute_is_registered(emitted):
|
||||
"An attribute added without a registry entry would be missing from the docs."
|
||||
unregistered = sorted(
|
||||
{
|
||||
f"{name} -> {key}"
|
||||
for name, kind, keys in emitted
|
||||
for key in keys
|
||||
if not reg.attribute_allowed(reg.span_for(name, kind), key)
|
||||
}
|
||||
)
|
||||
assert (
|
||||
not unregistered
|
||||
), "these span attributes are emitted but not registered: " + ", ".join(
|
||||
unregistered
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_every_registered_span_is_emitted(emitted):
|
||||
"""
|
||||
The direction nothing else catches: the docs must not describe a span that
|
||||
no longer exists.
|
||||
"""
|
||||
# By identity, not by name: a dynamic entry's own string never appears on
|
||||
# the wire, so comparing strings would be comparing the wrong things.
|
||||
resolved = {id(reg.span_for(name, kind)) for name, kind, _ in emitted}
|
||||
missing = sorted(str(span) for span in reg.SPANS if id(span) not in resolved)
|
||||
assert not missing, (
|
||||
f"these spans are documented but never emitted by the workload: {missing}. "
|
||||
"Either the instrumentation was removed, or exercise() no longer reaches it."
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_every_registered_attribute_is_emitted(emitted):
|
||||
"""
|
||||
Every registered attribute, optional or not, must actually be set at least
|
||||
once by the workload.
|
||||
|
||||
`optional` describes whether a reader should expect it on every span, not
|
||||
whether the code still sets it - so an attribute deleted from the code but
|
||||
left in the docs has to fail here even when it is marked optional. If a
|
||||
new attribute only appears in some rare case, extend exercise() to reach
|
||||
that case.
|
||||
"""
|
||||
by_entry = {}
|
||||
for name, kind, keys in emitted:
|
||||
entry = reg.span_for(name, kind)
|
||||
if entry is not None:
|
||||
by_entry.setdefault(id(entry), set()).update(keys)
|
||||
missing = []
|
||||
for span in reg.SPANS:
|
||||
emitted_keys = by_entry.get(id(span), set())
|
||||
for attribute in span.attributes:
|
||||
if attribute not in emitted_keys:
|
||||
missing.append(f"{span} -> {attribute}")
|
||||
assert not missing, (
|
||||
"these attributes are documented but never emitted by the workload: "
|
||||
+ ", ".join(sorted(missing))
|
||||
)
|
||||
|
||||
|
||||
def test_registry_has_no_duplicate_names():
|
||||
assert len(set(reg.SPANS)) == len(reg.SPANS)
|
||||
for span in reg.SPANS:
|
||||
assert len(set(span.attributes)) == len(
|
||||
span.attributes
|
||||
), f"{span} lists an attribute twice"
|
||||
|
||||
|
||||
def test_registry_entries_are_documented():
|
||||
"Every entry carries a description - the docs are generated from these."
|
||||
for span in reg.SPANS:
|
||||
assert span.description.strip(), f"{span} has no description"
|
||||
for attribute in span.attributes:
|
||||
assert attribute.description.strip(), f"{span} -> {attribute} has none"
|
||||
|
||||
|
||||
def test_registry_entries_are_usable_as_plain_strings():
|
||||
"The str subclassing is the whole reason call sites need no wrapper API."
|
||||
assert isinstance(reg.DB_QUERY, str)
|
||||
assert isinstance(reg.DB_NAMESPACE, str)
|
||||
assert reg.DB_QUERY == "db.query"
|
||||
assert reg.DB_NAMESPACE == "db.namespace"
|
||||
assert f"{reg.DB_QUERY}.execute" == "db.query.execute"
|
||||
|
||||
|
||||
def test_every_histogram_declares_bucket_boundaries():
|
||||
"""
|
||||
Every histogram must carry explicit boundaries, and only histograms may.
|
||||
|
||||
OpenTelemetry's default boundaries start at 5 and are meant for
|
||||
milliseconds, so a seconds-valued histogram that inherits them records
|
||||
everything into one bucket. This is a registry self-consistency check, not
|
||||
a check that the boundaries reached the SDK - for that see
|
||||
`test_histograms_spread_values_across_buckets` in test_telemetry_metrics.py.
|
||||
"""
|
||||
for metric in reg.METRICS:
|
||||
if metric.kind == reg.HISTOGRAM:
|
||||
assert metric.buckets, f"{metric} is a histogram with no boundaries"
|
||||
assert list(metric.buckets) == sorted(
|
||||
set(metric.buckets)
|
||||
), f"{metric} boundaries must be ascending and unique"
|
||||
assert metric.buckets[0] > 0, f"{metric} has a non-positive boundary"
|
||||
else:
|
||||
assert (
|
||||
metric.buckets is None
|
||||
), f"{metric} is a {metric.kind} and cannot have bucket boundaries"
|
||||
|
||||
|
||||
def test_dynamic_span_lookup():
|
||||
"""
|
||||
`dynamic=True` matching, which is how the request span resolves.
|
||||
|
||||
The last two assertions are the ones worth having: a dynamic entry must
|
||||
not swallow a span that does have a registered name, and must not match at
|
||||
all when the caller supplies no kind - otherwise every unregistered span
|
||||
in the suite would silently resolve to the request span and the
|
||||
emitted-but-not-registered direction would stop catching anything.
|
||||
"""
|
||||
assert reg.span_for("GET", SpanKind.SERVER) is reg.HTTP_REQUEST
|
||||
assert reg.span_for("POST /^/(?P<database>[^/]+)$", SpanKind.SERVER) is (
|
||||
reg.HTTP_REQUEST
|
||||
)
|
||||
assert reg.span_for("GET") is None
|
||||
assert reg.span_for("anything at all", SpanKind.INTERNAL) is None
|
||||
assert reg.span_for("db.query", SpanKind.SERVER) is reg.DB_QUERY
|
||||
|
||||
|
||||
def test_span_and_attribute_lookup():
|
||||
assert reg.span_for("db.query") is reg.DB_QUERY
|
||||
assert reg.span_for("datasette.startup") is reg.STARTUP
|
||||
assert reg.span_for("not.a.datasette.span") is None
|
||||
assert reg.attribute_allowed(reg.DB_QUERY, "db.namespace")
|
||||
assert not reg.attribute_allowed(reg.DB_QUERY, "db.namespace.extra")
|
||||
assert not reg.attribute_allowed(reg.DB_QUERY, "datasette.isolated_connection")
|
||||
assert not reg.attribute_allowed(None, "db.namespace")
|
||||
|
||||
|
||||
# --- Metric conformance ----------------------------------------------------
|
||||
|
||||
|
||||
@pytest_asyncio.fixture
|
||||
async def emitted_metrics(otel_metrics):
|
||||
"""
|
||||
Every (metric name, attribute key) pair produced by a broad workload,
|
||||
plus the raw set of metric names - the metric-side counterpart of the
|
||||
`emitted` span fixture above.
|
||||
|
||||
Metrics use DELTA temporality (see `otel_meter_provider` in datasette.telemetry_testing), and the
|
||||
function-scoped `otel_metrics` fixture drains any state left by an
|
||||
earlier test before yielding, so this collection is not polluted by
|
||||
other tests in the session - only by other *instances*, which is why the
|
||||
checks below key everything off attribute names rather than values.
|
||||
"""
|
||||
# The span workload already reaches every synchronous metric except the
|
||||
# interrupted counter: reads and writes drive db.client.operation.duration
|
||||
# and datasette.write.queue_wait, and both the suppressed-error probe and
|
||||
# the custom_time_limit interrupt raise through record_operation_duration,
|
||||
# setting error.type.
|
||||
ds = await exercise()
|
||||
|
||||
# datasette.sql.queries.interrupted counts only queries that exceed the
|
||||
# *configured* limit - a caller opting into a deliberately short budget
|
||||
# via custom_time_limit (as exercise() does) is excluded by design. So a
|
||||
# second instance whose configured limit is tiny provides the real thing.
|
||||
slow_name = _unique("registry_metrics_slow")
|
||||
slow = Datasette(memory=True, settings={"sql_time_limit_ms": 5})
|
||||
slow.add_memory_database(slow_name)
|
||||
await slow.invoke_startup()
|
||||
slow_db = slow.get_database(slow_name)
|
||||
with pytest.raises(QueryInterrupted):
|
||||
await slow_db.execute(
|
||||
"with recursive c(x) as (select 0 union all select x+1 from c) "
|
||||
"select * from c"
|
||||
)
|
||||
|
||||
# Collect while both instances are still registered, so the observable
|
||||
# gauges - which observe live instances at collection time - report.
|
||||
otel_metrics.collect()
|
||||
snapshot = otel_metrics.snapshot
|
||||
assert snapshot, "no metrics captured - the fixture is not exercising anything"
|
||||
pairs = set()
|
||||
for metric_name, points in snapshot.items():
|
||||
for point in points:
|
||||
for key in point.attributes or {}:
|
||||
pairs.add((metric_name, key))
|
||||
ds.close()
|
||||
slow.close()
|
||||
return {"names": set(snapshot), "pairs": pairs, "collector": otel_metrics}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_metrics_conform_to_the_registry(emitted_metrics):
|
||||
"""
|
||||
Emitted-but-unregistered, via the plugin kit's helper - consumed here
|
||||
exactly the way a plugin's suite would. Beyond names and attribute keys,
|
||||
this also asserts each instrument was created as the kind and unit its
|
||||
registry entry declares, and that `datasette.operation` only ever takes
|
||||
its declared enum values.
|
||||
"""
|
||||
assert_metrics_conform(
|
||||
reg.METRICS, emitted_metrics["collector"], scope_name="datasette"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_every_registered_metric_is_emitted(emitted_metrics):
|
||||
"Registered-but-never-collected, via the plugin kit's helper."
|
||||
assert_metrics_covered(
|
||||
reg.METRICS, emitted_metrics["collector"], scope_name="datasette"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_every_registered_metric_attribute_is_emitted(emitted_metrics):
|
||||
"""
|
||||
The direction nothing else catches: the docs must not describe a metric
|
||||
attribute that no longer exists.
|
||||
|
||||
Unlike the span-side attribute check, this does not skip `optional`
|
||||
attributes. The only optional metric attribute is `error.type` on
|
||||
`db.client.operation.duration`, and the workload reaches it from two
|
||||
independent directions: the suppressed-error probe and the
|
||||
custom_time_limit interrupt in `exercise()`, both of which raise through
|
||||
`record_operation_duration`. So it is checked like any other attribute
|
||||
rather than exempted; marking something optional here would opt it out of
|
||||
verification entirely.
|
||||
|
||||
Gauges with no registered attributes (`datasette.sql.threads.limit` and
|
||||
`.queue_depth`) fall out correctly with no special case: their
|
||||
`metric.attributes` is empty, so the inner loop makes no assertion.
|
||||
"""
|
||||
emitted_keys_by_metric = {}
|
||||
for metric_name, key in emitted_metrics["pairs"]:
|
||||
emitted_keys_by_metric.setdefault(metric_name, set()).add(key)
|
||||
|
||||
missing = []
|
||||
for metric in reg.METRICS:
|
||||
if str(metric) not in emitted_metrics["names"]:
|
||||
# Not emitted at all - already reported by
|
||||
# test_every_registered_metric_is_emitted; do not double-report.
|
||||
continue
|
||||
emitted_keys = emitted_keys_by_metric.get(str(metric), set())
|
||||
for attribute in metric.attributes:
|
||||
if attribute not in emitted_keys:
|
||||
missing.append(f"{metric} -> {attribute}")
|
||||
assert not missing, (
|
||||
"these metric attributes are documented but never emitted by the "
|
||||
"test workload: " + ", ".join(sorted(missing))
|
||||
)
|
||||
|
||||
|
||||
def test_prefix_span_lookup():
|
||||
"""
|
||||
`prefix=True` matching, exercised directly.
|
||||
|
||||
Core registers no prefix spans - the flag exists for plugin registries
|
||||
(e.g. a `chat {model}` span family) - so without this the branch in
|
||||
`span_for()` would be untested code the conformance tests never reach.
|
||||
"""
|
||||
hook = reg.SpanName("myplugin.hook.", "A hypothetical span family", prefix=True)
|
||||
spans = reg.SPANS + (hook,)
|
||||
assert reg.span_for("myplugin.hook.render_cell", spans=spans) is hook
|
||||
assert reg.span_for("myplugin.hook.anything", spans=spans) is hook
|
||||
assert reg.span_for("myplugin.hookish", spans=spans) is None
|
||||
assert reg.span_for("db.query", spans=spans) is reg.DB_QUERY
|
||||
|
||||
|
||||
def test_exact_match_wins_over_prefix():
|
||||
"A prefix family can never shadow a span with a registered exact name."
|
||||
family = reg.SpanName("db.", "Greedy prefix", prefix=True)
|
||||
spans = (family,) + reg.SPANS
|
||||
assert reg.span_for("db.query", spans=spans) is reg.DB_QUERY
|
||||
assert reg.span_for("db.anything-else", spans=spans) is family
|
||||
|
||||
|
||||
def test_attribute_values_enum_enforced():
|
||||
outcome = reg.Attribute("myplugin.outcome", "Enum.", values={"ok", "error"})
|
||||
open_attr = reg.Attribute("myplugin.note", "Open value set.")
|
||||
span = reg.SpanName("myplugin.job", "Test span", (outcome, open_attr))
|
||||
assert reg.attribute_value_allowed(span, "myplugin.outcome", "ok")
|
||||
assert not reg.attribute_value_allowed(span, "myplugin.outcome", "surprise")
|
||||
assert reg.attribute_value_allowed(span, "myplugin.note", "anything at all")
|
||||
assert not reg.attribute_value_allowed(span, "not.registered", "x")
|
||||
assert not reg.attribute_value_allowed(None, "myplugin.outcome", "ok")
|
||||
|
|
@ -1,331 +0,0 @@
|
|||
"""
|
||||
The plugin telemetry kit (`datasette.telemetry_testing` plus the public
|
||||
registry classes), exercised the way a third-party plugin would use it: a
|
||||
toy plugin registry, a toy tracer scope, and the kit's own fixtures and
|
||||
conformance helpers.
|
||||
"""
|
||||
|
||||
import pytest
|
||||
|
||||
pytest.importorskip("opentelemetry.sdk")
|
||||
|
||||
from opentelemetry import trace as otel_trace
|
||||
|
||||
from datasette import telemetry_registry as reg
|
||||
from datasette.telemetry import linked_root_span_kwargs
|
||||
from datasette.telemetry_testing import (
|
||||
assert_package_never_imports_sdk,
|
||||
assert_spans_covered,
|
||||
assert_spans_conform,
|
||||
)
|
||||
|
||||
SCOPE = "toyplugin"
|
||||
|
||||
OUTCOME = reg.Attribute(
|
||||
"toyplugin.outcome", "How the job ended.", values={"ok", "error"}
|
||||
)
|
||||
JOB_NAME = reg.Attribute("toyplugin.job", "The job's registered name.")
|
||||
JOB = reg.SpanName("toyplugin.job.run", "One job execution.", (OUTCOME, JOB_NAME))
|
||||
CHAT = reg.SpanName(
|
||||
"toyplugin.chat ", "One model call, named `toyplugin.chat {model}`.", prefix=True
|
||||
)
|
||||
TOY_SPANS = (JOB, CHAT)
|
||||
|
||||
toy_tracer = otel_trace.get_tracer(SCOPE, "0.1")
|
||||
|
||||
|
||||
def _toy_spans(otel_spans):
|
||||
return [
|
||||
span
|
||||
for span in otel_spans.get_finished_spans()
|
||||
if span.instrumentation_scope and span.instrumentation_scope.name == SCOPE
|
||||
]
|
||||
|
||||
|
||||
def _run_workload():
|
||||
with toy_tracer.start_as_current_span(JOB) as span:
|
||||
span.set_attribute(OUTCOME, "ok")
|
||||
span.set_attribute(JOB_NAME, "nightly")
|
||||
with toy_tracer.start_as_current_span("toyplugin.chat gpt-5"):
|
||||
pass
|
||||
|
||||
|
||||
def test_conformance_passes_for_a_conforming_workload(otel_spans):
|
||||
_run_workload()
|
||||
finished = otel_spans.get_finished_spans()
|
||||
assert_spans_conform(TOY_SPANS, finished, scope_name=SCOPE)
|
||||
# Coverage direction needs prefix families seen too - the chat span
|
||||
# resolves to the CHAT entry despite its variable suffix.
|
||||
assert_spans_covered(TOY_SPANS, finished, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_conformance_catches_an_unregistered_span(otel_spans):
|
||||
with toy_tracer.start_as_current_span("toyplugin.surprise"):
|
||||
pass
|
||||
with pytest.raises(AssertionError, match="unregistered span"):
|
||||
assert_spans_conform(
|
||||
TOY_SPANS, otel_spans.get_finished_spans(), scope_name=SCOPE
|
||||
)
|
||||
|
||||
|
||||
def test_conformance_catches_an_unregistered_attribute(otel_spans):
|
||||
with toy_tracer.start_as_current_span(JOB) as span:
|
||||
span.set_attribute("toyplugin.stealth", 1)
|
||||
with pytest.raises(AssertionError, match="unregistered attribute"):
|
||||
assert_spans_conform(
|
||||
TOY_SPANS, otel_spans.get_finished_spans(), scope_name=SCOPE
|
||||
)
|
||||
|
||||
|
||||
def test_conformance_enforces_declared_enums(otel_spans):
|
||||
with toy_tracer.start_as_current_span(JOB) as span:
|
||||
span.set_attribute(OUTCOME, "surprise")
|
||||
with pytest.raises(AssertionError, match="not in the declared enum"):
|
||||
assert_spans_conform(
|
||||
TOY_SPANS, otel_spans.get_finished_spans(), scope_name=SCOPE
|
||||
)
|
||||
|
||||
|
||||
def test_coverage_catches_a_never_emitted_span(otel_spans):
|
||||
with toy_tracer.start_as_current_span(JOB) as span:
|
||||
span.set_attribute(OUTCOME, "ok")
|
||||
span.set_attribute(JOB_NAME, "nightly")
|
||||
# CHAT never emitted
|
||||
with pytest.raises(AssertionError, match="never emitted"):
|
||||
assert_spans_covered(
|
||||
TOY_SPANS, otel_spans.get_finished_spans(), scope_name=SCOPE
|
||||
)
|
||||
|
||||
|
||||
def test_scope_filter_ignores_other_scopes(otel_spans):
|
||||
# Core's own spans are in the exporter too; a plugin's conformance run
|
||||
# must not fail because of them.
|
||||
other = otel_trace.get_tracer("someone-else", "1.0")
|
||||
with other.start_as_current_span("not.in.the.toy.registry"):
|
||||
pass
|
||||
_run_workload()
|
||||
assert_spans_conform(TOY_SPANS, otel_spans.get_finished_spans(), scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_linked_root_span_kwargs_links_without_parenting(otel_spans):
|
||||
with toy_tracer.start_as_current_span("toyplugin.cause") as cause:
|
||||
cause_context = cause.get_span_context()
|
||||
kwargs = linked_root_span_kwargs()
|
||||
with toy_tracer.start_as_current_span("toyplugin.effect", **kwargs):
|
||||
pass
|
||||
effect = [
|
||||
span for span in _toy_spans(otel_spans) if span.name == "toyplugin.effect"
|
||||
][0]
|
||||
assert effect.parent is None, "must be a root, not a child"
|
||||
assert effect.context.trace_id != cause_context.trace_id
|
||||
assert len(effect.links) == 1
|
||||
assert effect.links[0].context.span_id == cause_context.span_id
|
||||
|
||||
|
||||
def test_linked_root_span_kwargs_with_no_current_span(otel_spans):
|
||||
kwargs = linked_root_span_kwargs()
|
||||
assert kwargs["links"] == []
|
||||
with toy_tracer.start_as_current_span("toyplugin.orphanless", **kwargs):
|
||||
pass
|
||||
span = _toy_spans(otel_spans)[0]
|
||||
assert span.parent is None
|
||||
assert span.links == ()
|
||||
|
||||
|
||||
def test_kit_module_itself_never_imports_the_sdk():
|
||||
"""
|
||||
The kit imports the SDK lazily, so a plugin importing it at module
|
||||
level does not violate the api-only dependency rule.
|
||||
|
||||
conftest.py's pytest_collection_modifyitems() moves this test to the
|
||||
front of the run by name - if you rename it, rename it there too. Like
|
||||
every subprocess-spawning test in this suite, running it late crashes
|
||||
the interpreter on macOS/CPython 3.13 (SIGBUS in fork+exec once the
|
||||
process holds enough threads) - see the comment there.
|
||||
"""
|
||||
assert_package_never_imports_sdk("datasette.telemetry_testing")
|
||||
|
||||
|
||||
# --- Metric conformance helpers --------------------------------------------
|
||||
|
||||
import itertools
|
||||
|
||||
from opentelemetry import metrics as otel_metrics_api
|
||||
|
||||
from datasette.telemetry_testing import (
|
||||
assert_metrics_conform,
|
||||
assert_metrics_covered,
|
||||
)
|
||||
|
||||
toy_meter = otel_metrics_api.get_meter(SCOPE, "0.1")
|
||||
|
||||
# Instrument names must be unique per meter for the SDK, so each test mints
|
||||
# its own via this counter rather than re-registering one name.
|
||||
_metric_ids = itertools.count()
|
||||
|
||||
|
||||
def _toy_metric_registry(name, kind="Counter", unit="{job}", attributes=None):
|
||||
return (
|
||||
reg.MetricName(
|
||||
name,
|
||||
kind,
|
||||
unit,
|
||||
"A toy metric.",
|
||||
attributes if attributes is not None else (OUTCOME,),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def test_metrics_conform_passes_and_covers(otel_metrics):
|
||||
name = f"toyplugin.jobs.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name)
|
||||
counter = toy_meter.create_counter(name, unit="{job}", description="Jobs run")
|
||||
counter.add(1, {OUTCOME: "ok"})
|
||||
otel_metrics.collect()
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
assert_metrics_covered(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_metrics_conform_catches_unregistered_metric(otel_metrics):
|
||||
name = f"toyplugin.stealth.{next(_metric_ids)}"
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1)
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="unregistered metric"):
|
||||
assert_metrics_conform((), otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_metrics_conform_catches_kind_mismatch(otel_metrics):
|
||||
name = f"toyplugin.kindclash.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name, kind="Histogram", unit="{job}")
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {OUTCOME: "ok"})
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="registry declares Histogram"):
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_metrics_conform_catches_unit_mismatch(otel_metrics):
|
||||
name = f"toyplugin.unitclash.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name, unit="s")
|
||||
counter = toy_meter.create_counter(name, unit="ms")
|
||||
counter.add(1, {OUTCOME: "ok"})
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="unit"):
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_metrics_conform_catches_unregistered_attribute(otel_metrics):
|
||||
name = f"toyplugin.attrclash.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name)
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {"toyplugin.stealth": "x"})
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="unregistered attribute"):
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_metrics_conform_enforces_declared_enums(otel_metrics):
|
||||
name = f"toyplugin.enumclash.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name)
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {OUTCOME: "surprise"})
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="not in the declared enum"):
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_metrics_covered_catches_never_collected(otel_metrics):
|
||||
registered_but_never_created = _toy_metric_registry(
|
||||
f"toyplugin.ghost.{next(_metric_ids)}"
|
||||
)
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="never collected"):
|
||||
assert_metrics_covered(
|
||||
registered_but_never_created, otel_metrics, scope_name=SCOPE
|
||||
)
|
||||
|
||||
|
||||
def test_metrics_covered_skips_optional_attributes(otel_metrics):
|
||||
name = f"toyplugin.optattr.{next(_metric_ids)}"
|
||||
error_type = reg.Attribute("toyplugin.error", "Only on failure.", optional=True)
|
||||
registry = _toy_metric_registry(name, attributes=(OUTCOME, error_type))
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {OUTCOME: "ok"}) # no error attribute - and that is fine
|
||||
otel_metrics.collect()
|
||||
assert_metrics_covered(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_metrics_scope_filter_ignores_other_scopes(otel_metrics):
|
||||
# Core's own metrics are in the reader too; a plugin's conformance run
|
||||
# must not fail because of them.
|
||||
name = f"toyplugin.scoped.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name)
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {OUTCOME: "ok"})
|
||||
other_meter = otel_metrics_api.get_meter("someone-else-metrics", "1.0")
|
||||
stranger = other_meter.create_counter(f"stranger.{next(_metric_ids)}", unit="x")
|
||||
stranger.add(1)
|
||||
otel_metrics.collect()
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
# --- UpDownCounter kind + privacy walk --------------------------------------
|
||||
|
||||
from datasette.telemetry_testing import assert_no_forbidden_values
|
||||
|
||||
|
||||
def test_updown_counter_kind_passes(otel_metrics):
|
||||
name = f"toyplugin.active.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name, kind=reg.UPDOWN_COUNTER, unit="{turn}")
|
||||
updown = toy_meter.create_up_down_counter(name, unit="{turn}")
|
||||
updown.add(1, {OUTCOME: "ok"})
|
||||
otel_metrics.collect()
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_counter_registered_as_updown_fails_on_monotonicity(otel_metrics):
|
||||
name = f"toyplugin.monoclash.{next(_metric_ids)}"
|
||||
registry = _toy_metric_registry(name, kind=reg.UPDOWN_COUNTER, unit="{job}")
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {OUTCOME: "ok"})
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="is_monotonic"):
|
||||
assert_metrics_conform(registry, otel_metrics, scope_name=SCOPE)
|
||||
|
||||
|
||||
def test_forbidden_values_walk_catches_a_leak(otel_spans, otel_metrics):
|
||||
secret = "sentinel-token-xyzzy"
|
||||
with toy_tracer.start_as_current_span(JOB) as span:
|
||||
span.set_attribute(OUTCOME, "ok")
|
||||
span.set_attribute(JOB_NAME, f"job for {secret}")
|
||||
with pytest.raises(AssertionError, match="sentinel-token-xyzzy"):
|
||||
assert_no_forbidden_values(
|
||||
{secret},
|
||||
finished_spans=otel_spans.get_finished_spans(),
|
||||
scope_name=SCOPE,
|
||||
)
|
||||
|
||||
|
||||
def test_forbidden_values_walk_passes_a_clean_workload(otel_spans, otel_metrics):
|
||||
_run_workload()
|
||||
name = f"toyplugin.clean.{next(_metric_ids)}"
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {OUTCOME: "ok"})
|
||||
otel_metrics.collect()
|
||||
assert_no_forbidden_values(
|
||||
{"sentinel-token-xyzzy", "alice@example.com", ""},
|
||||
finished_spans=otel_spans.get_finished_spans(),
|
||||
collector=otel_metrics,
|
||||
scope_name=SCOPE,
|
||||
)
|
||||
|
||||
|
||||
def test_forbidden_values_walk_checks_metric_attributes(otel_metrics):
|
||||
secret = "leaky-metric-value"
|
||||
name = f"toyplugin.leak.{next(_metric_ids)}"
|
||||
counter = toy_meter.create_counter(name, unit="{job}")
|
||||
counter.add(1, {"toyplugin.note": secret})
|
||||
otel_metrics.collect()
|
||||
with pytest.raises(AssertionError, match="leaky-metric-value"):
|
||||
assert_no_forbidden_values({secret}, collector=otel_metrics, scope_name=SCOPE)
|
||||
Loading…
Add table
Add a link
Reference in a new issue