Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion ddtrace/contrib/_events/web_framework.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,16 @@
from ddtrace.contrib._events.http import HttpRequestBaseEvent
from ddtrace.ext import SpanKind
from ddtrace.ext import SpanTypes
from ddtrace.internal.core.event_names import WEB_REQUEST as WEB_REQUEST_EVENT
from ddtrace.internal.core.event_names import WEB_REQUEST_STARTING as WEB_REQUEST_STARTING_EVENT
from ddtrace.internal.core.events import event_field
from ddtrace.internal.schema import SpanDirection
from ddtrace.internal.schema import schematize_url_operation


class WebFrameworkEvents(str, Enum):
WEB_REQUEST = "web.request"
WEB_REQUEST = WEB_REQUEST_EVENT
WEB_REQUEST_STARTING = WEB_REQUEST_STARTING_EVENT


@dataclass
Expand Down
4 changes: 4 additions & 0 deletions ddtrace/contrib/internal/flask/patch.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from werkzeug.exceptions import NotFound

from ddtrace.contrib import trace_utils
from ddtrace.contrib._events.web_framework import WebFrameworkEvents
from ddtrace.ext import SpanTypes
from ddtrace.internal import core
from ddtrace.internal.constants import COMPONENT
Expand Down Expand Up @@ -391,6 +392,9 @@ def unpatch():

def patched_wsgi_app(wrapped, instance, args, kwargs):
environ, start_response = args
core.dispatch(
WebFrameworkEvents.WEB_REQUEST_STARTING.value, (environ.get("REQUEST_METHOD"), environ.get("PATH_INFO"))
)
# Registration is gated on asm_config, not tracing — keep this above the tracing short-circuit.
_collect_routes_once(instance, environ.get("SCRIPT_NAME") or "")
if not is_tracing_enabled():
Expand Down
45 changes: 45 additions & 0 deletions ddtrace/internal/core/crashtracking.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
from ddtrace.internal.compat import ensure_text
from ddtrace.internal.logger import get_logger
from ddtrace.internal.runtime import get_runtime_id
from ddtrace.internal.runtime import on_runtime_id_change
from ddtrace.internal.settings import env
from ddtrace.internal.settings._agent import config as agent_config
from ddtrace.internal.settings.crashtracker import config as crashtracker_config
Expand All @@ -33,6 +34,7 @@
from ddtrace.internal.native._native import StacktraceCollection
from ddtrace.internal.native._native import crashtracker_init
from ddtrace.internal.native._native import crashtracker_on_fork
from ddtrace.internal.native._native import crashtracker_reconfigure
from ddtrace.internal.native._native import crashtracker_report_unhandled_exception
from ddtrace.internal.native._native import crashtracker_status

Expand All @@ -41,6 +43,26 @@
is_available = False


# Runtime ID callbacks only receive the new ID, so keep the start() tags here
# to rebuild crashtracker metadata on explicit identity refreshes.
_identity_refresh_additional_tags: Optional[dict[str, str]] = None
_identity_refresh_fork_generation = forksafe.get_generation()


def _on_identity_refresh(_new_runtime_id: str) -> None:
global _identity_refresh_fork_generation

fork_generation = forksafe.get_generation()
if fork_generation != _identity_refresh_fork_generation:
# Runtime ID callbacks also run from the forksafe hook, before crashtracker's
# own fork handler. Skip that notification so the native crashtracker is
# reset only by crashtracker_on_fork; explicit later identity refreshes in
# the child still reconfigure below.
_identity_refresh_fork_generation = fork_generation
return
_reconfigure_for_identity_refresh(_identity_refresh_additional_tags)


def _get_tags(additional_tags: Optional[dict[str, str]]) -> dict[str, str]:
tags = {
"language": "python",
Expand Down Expand Up @@ -211,12 +233,32 @@ def is_started() -> bool:
return crashtracker_status() == CrashtrackerStatus.Initialized


def _reconfigure_for_identity_refresh(additional_tags: Optional[dict[str, str]]) -> None:
# Native reconfigure is best-effort during identity refresh: if it fails,
# do not break runtime.refresh_identity() or later runtime ID callbacks.
try:
if not is_started():
return

config, receiver_config, metadata = _get_args(additional_tags)
if config is None or receiver_config is None or metadata is None:
log.error("Failed to reconfigure crashtracker after identity refresh: failed to construct configuration")
return
crashtracker_reconfigure(config, receiver_config, metadata)
except Exception:
log.debug("Failed to reconfigure crashtracker after identity refresh", exc_info=True)


def start(additional_tags: Optional[dict[str, str]] = None) -> bool:
global _identity_refresh_additional_tags, _identity_refresh_fork_generation

if not is_available:
return False
if not crashtracker_config.enabled:
return False

additional_tags = dict(additional_tags) if additional_tags is not None else None

try:
config, receiver_config, metadata = _get_args(additional_tags)
if config is None or receiver_config is None or metadata is None:
Expand Down Expand Up @@ -256,6 +298,9 @@ def crashtracker_fork_handler():
crashtracker_on_fork(config, receiver_config, metadata)

forksafe.register(crashtracker_fork_handler)
_identity_refresh_additional_tags = additional_tags
_identity_refresh_fork_generation = forksafe.get_generation()
on_runtime_id_change(_on_identity_refresh)
except Exception:
log.exception("Failed to start crashtracker")
return False
Expand Down
9 changes: 9 additions & 0 deletions ddtrace/internal/core/event_names.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
"""Shared names for events dispatched through ddtrace.internal.core."""

WEB_REQUEST = "web.request"
WEB_REQUEST_STARTING = "web.request.starting"

__all__ = [
"WEB_REQUEST",
"WEB_REQUEST_STARTING",
]
3 changes: 3 additions & 0 deletions ddtrace/internal/native/_native.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,9 @@ def crashtracker_init(
def crashtracker_on_fork(
config: CrashtrackerConfiguration, receiver_config: CrashtrackerReceiverConfig, metadata: CrashtrackerMetadata
) -> None: ...
def crashtracker_reconfigure(
config: CrashtrackerConfiguration, receiver_config: CrashtrackerReceiverConfig, metadata: CrashtrackerMetadata
) -> None: ...
def crashtracker_status() -> CrashtrackerStatus: ...
def crashtracker_receiver() -> None: ...
def crashtracker_report_unhandled_exception(
Expand Down
13 changes: 13 additions & 0 deletions src/native/crashtracker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -320,6 +320,19 @@ pub fn crashtracker_on_fork<'py>(
libdd_crashtracker::on_fork(inner_config, inner_receiver_config, inner_metadata)
}

#[pyfunction(name = "crashtracker_reconfigure")]
pub fn crashtracker_reconfigure<'py>(
mut config: PyRefMut<'py, CrashtrackerConfigurationPy>,
mut receiver_config: PyRefMut<'py, CrashtrackerReceiverConfigPy>,
mut metadata: PyRefMut<'py, CrashtrackerMetadataPy>,
) -> anyhow::Result<()> {
let inner_config = (*config).take_inner_or_err()?;
let inner_receiver_config = (*receiver_config).take_inner_or_err()?;
let inner_metadata = (*metadata).take_inner_or_err()?;

libdd_crashtracker::reconfigure(inner_config, inner_receiver_config, inner_metadata)
}

#[pyfunction(name = "crashtracker_status")]
pub fn crashtracker_status() -> anyhow::Result<CrashtrackerStatus> {
CrashtrackerStatus::try_from(CRASHTRACKER_STATUS.load(Ordering::SeqCst))
Expand Down
1 change: 1 addition & 0 deletions src/native/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ fn _native(m: &Bound<'_, PyModule>) -> PyResult<()> {
m.add_class::<crashtracker::CrashtrackerStatus>()?;
m.add_function(wrap_pyfunction!(crashtracker::crashtracker_init, m)?)?;
m.add_function(wrap_pyfunction!(crashtracker::crashtracker_on_fork, m)?)?;
m.add_function(wrap_pyfunction!(crashtracker::crashtracker_reconfigure, m)?)?;
m.add_function(wrap_pyfunction!(crashtracker::crashtracker_status, m)?)?;
m.add_function(wrap_pyfunction!(crashtracker::crashtracker_receiver, m)?)?;
m.add_function(wrap_pyfunction!(
Expand Down
71 changes: 71 additions & 0 deletions tests/contrib/flask/test_microvm_identity_refresh.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
import mock

from ddtrace.contrib._events.web_framework import WebFrameworkEvents
from ddtrace.contrib.internal.flask.patch import patched_wsgi_app
from ddtrace.internal import core

from . import BaseFlaskTestCase


REQUEST_STARTING_PATH = "/web-request-starting"


class FlaskMicrovmIdentityRefreshTestCase(BaseFlaskTestCase):
"""patched_wsgi_app() must dispatch every request's method/path before tracing starts.

The matching logic itself is tested in tests/tracer/runtime/test_runtime_id.py.
"""

def test_microvm_run_hook_request(self):
"""No route is registered at the hook path: wsgi_app() runs before routing, so it
must fire even on a 404 -- the stronger, more general form of this check (whether the
route matches doesn't change what gets dispatched).
"""
with mock.patch("ddtrace.contrib.internal.flask.patch.core.dispatch", wraps=core.dispatch) as m:
res = self.client.post(REQUEST_STARTING_PATH)

self.assertEqual(res.status_code, 404)
m.assert_any_call(WebFrameworkEvents.WEB_REQUEST_STARTING.value, ("POST", REQUEST_STARTING_PATH))

def test_other_request(self):
@self.app.route("/")
def index():
return "ok", 200

with mock.patch("ddtrace.contrib.internal.flask.patch.core.dispatch", wraps=core.dispatch) as m:
res = self.client.get("/")

self.assertEqual(res.status_code, 200)
m.assert_any_call(WebFrameworkEvents.WEB_REQUEST_STARTING.value, ("GET", "/"))

def test_pre_request_event_dispatches_before_wsgi_middleware(self):
events = []
environ = {"REQUEST_METHOD": "POST", "PATH_INFO": REQUEST_STARTING_PATH, "SCRIPT_NAME": ""}

def start_response(status, headers, exc_info=None):
pass

def wrapped(environ, start_response):
return []

def dispatch(name, args):
if name == WebFrameworkEvents.WEB_REQUEST_STARTING.value:
events.append("starting")

class WSGIMiddleware:
def __init__(self, app, tracer, integration_config):
pass

def __call__(self, environ, start_response):
events.append("middleware")
return []

with (
mock.patch("ddtrace.contrib.internal.flask.patch.core.dispatch", side_effect=dispatch),
mock.patch("ddtrace.contrib.internal.flask.patch._collect_routes_once"),
mock.patch("ddtrace.contrib.internal.flask.patch.is_tracing_enabled", return_value=True),
mock.patch("ddtrace.contrib.internal.flask.patch._FlaskWSGIMiddleware", WSGIMiddleware),
):
patched_wsgi_app(wrapped, self.app, (environ, start_response), {})

assert events == ["starting", "middleware"]
Loading