mirror of
https://github.com/element-hq/synapse.git
synced 2026-08-25 03:09:55 +00:00
Port LoggingContext to a Rust pyclass
Move `LoggingContext` from a pure-Python class to a native `#[pyclass]` in `rust/src/logcontext.rs`, re-exported from `synapse.logging.context`. The attribute surface, methods, error-message wording and abuse-detection behaviour are all preserved, so callers — and the Python subclass `BackgroundProcessLoggingContext` — are unaffected. `set_current_context` stays in Python for now: it still calls `ctx.start`/ `ctx.stop` via normal attribute dispatch, so subclass overrides continue to work. The native switch primitive (with in-Rust `getrusage`) is a later increment. Notes on the port: - Construction is split between `__new__` (allocates a blank instance) and `__init__` (does the real work), mirroring a pure-Python class. This is what lets `BackgroundProcessLoggingContext`, which composes a name and then calls `super().__init__(name=..., server_name=...)`, keep working unchanged. - Methods that both mutate `self` and call back into Python (`start`, `stop`, `__enter__`, `__exit__`) take an owned `Bound<Self>` and only hold a `borrow_mut()` across the field write itself, never across a Python call, so a re-entrant log record (which resolves the current context and reads this object's getters) cannot hit an "already borrowed" error. - `logcontext_error`, `set_current_context`, `logger` and `logcontext_debug_logger` are re-resolved through the module on every call so that test patches take effect. - `__traverse__`/`__clear__` cover `previous_context`/`parent_context`/ `request`/`scope` for the cyclic GC (scope <-> context is a real cycle). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011LnJSXQ86AAWNtihR4C5PH
This commit is contained in:
co-authored by
Claude Fable 5
parent
fa91dc5e3c
commit
ad59e71e2a
@@ -46,8 +46,13 @@
|
||||
|
||||
use std::{cell::RefCell, future::Future};
|
||||
|
||||
use log::{debug, log_enabled, Level};
|
||||
use once_cell::sync::OnceCell;
|
||||
use pyo3::call::PyCallArgs;
|
||||
use pyo3::exceptions::{PyAssertionError, PyValueError};
|
||||
use pyo3::prelude::*;
|
||||
use pyo3::types::{PyDict, PyTuple};
|
||||
use pyo3::{PyTraverseError, PyVisit};
|
||||
|
||||
/// The Python sentinel logcontext (`synapse.logging.context.SENTINEL_CONTEXT`).
|
||||
///
|
||||
@@ -56,6 +61,18 @@ use pyo3::prelude::*;
|
||||
/// must not import `synapse.logging.context`; see [`crate::deferred`]).
|
||||
static SENTINEL: OnceCell<Py<PyAny>> = OnceCell::new();
|
||||
|
||||
/// Name of the opt-in logger for logcontext switch tracing.
|
||||
///
|
||||
/// This is the single source of truth for the logger name: it is used as the
|
||||
/// `debug!` `target:` for the switch traces emitted below, and it is exported to
|
||||
/// Python via [`register_module`] so that `synapse.logging.context` builds
|
||||
/// *exactly* this logger (see `logcontext_debug_logger` there). Keeping one
|
||||
/// constant stops the Rust `target:` and the Python `getLogger` name from
|
||||
/// drifting apart. The messages only surface when this logger is explicitly
|
||||
/// configured — see `ExplicitlyConfiguredLogger` on the Python side, whose
|
||||
/// `isEnabledFor` pyo3-log honours, so the no-inherit opt-in works from Rust too.
|
||||
pub const DEBUG_LOGGER_NAME: &str = "synapse.logging.context.debug";
|
||||
|
||||
thread_local! {
|
||||
/// The current logcontext for this OS thread. `None` means "no context set on
|
||||
/// this thread", which is reported as the sentinel — matching the old
|
||||
@@ -235,6 +252,538 @@ impl ContextResourceUsage {
|
||||
}
|
||||
}
|
||||
|
||||
/// Import `synapse.logging.context`.
|
||||
///
|
||||
/// After the module is first imported this is just a `sys.modules` lookup, so it
|
||||
/// is cheap enough to call on the switch path. We deliberately re-resolve
|
||||
/// `logcontext_error` / `set_current_context` / `logger` through the module on
|
||||
/// every call (rather than caching the callables) so that test patches of those
|
||||
/// module-level names take effect, and so we never import the module at Rust
|
||||
/// module-registration time (which would be a circular import).
|
||||
fn context_module(py: Python<'_>) -> PyResult<Bound<'_, PyModule>> {
|
||||
py.import("synapse.logging.context")
|
||||
}
|
||||
|
||||
/// Call the (possibly test-patched) module-level `logcontext_error(msg)`.
|
||||
fn logcontext_error(module: &Bound<'_, PyModule>, msg: String) -> PyResult<()> {
|
||||
module.getattr("logcontext_error")?.call1((msg,))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// `threading.get_ident`, cached (it never changes and this is on a hot path).
|
||||
static GET_THREAD_ID: OnceCell<Py<PyAny>> = OnceCell::new();
|
||||
|
||||
/// The current OS thread id, matching Python's `threading.get_ident()`.
|
||||
fn get_thread_id(py: Python<'_>) -> PyResult<u64> {
|
||||
let get_ident = GET_THREAD_ID.get_or_try_init(|| -> PyResult<Py<PyAny>> {
|
||||
Ok(py.import("threading")?.getattr("get_ident")?.unbind())
|
||||
})?;
|
||||
get_ident.bind(py).call0()?.extract()
|
||||
}
|
||||
|
||||
/// Normalise a Python value to `None` if it is `None`, else `Some`.
|
||||
fn none_to_option(obj: Bound<'_, PyAny>) -> Option<Py<PyAny>> {
|
||||
if obj.is_none() {
|
||||
None
|
||||
} else {
|
||||
Some(obj.unbind())
|
||||
}
|
||||
}
|
||||
|
||||
/// Python truthiness of an optional slot: `None` is falsy, otherwise the object's
|
||||
/// own `bool()`. Used for the "is this context active?" (`usage_start`) and "did
|
||||
/// we get a real rusage?" checks, matching Python's `if self.usage_start:` /
|
||||
/// `if not rusage:`.
|
||||
fn is_truthy(py: Python<'_>, slot: &Option<Py<PyAny>>) -> PyResult<bool> {
|
||||
match slot {
|
||||
Some(obj) => obj.bind(py).is_truthy(),
|
||||
None => Ok(false),
|
||||
}
|
||||
}
|
||||
|
||||
/// Propagate a usage update to the parent context, if there is a (truthy) one.
|
||||
///
|
||||
/// Dispatched via `call_method1` so subclass overrides are respected. The
|
||||
/// truthiness guard matches Python's `if self.parent_context:` and is
|
||||
/// load-bearing: the sentinel is falsy *and* implements no `add_*` methods, so a
|
||||
/// bare `is_some()` check would call a nonexistent method on it.
|
||||
fn forward_to_parent<'py>(
|
||||
parent: &Option<Py<PyAny>>,
|
||||
py: Python<'py>,
|
||||
method: &str,
|
||||
args: impl PyCallArgs<'py>,
|
||||
) -> PyResult<()> {
|
||||
if let Some(parent) = parent {
|
||||
let parent = parent.bind(py);
|
||||
if parent.is_truthy()? {
|
||||
parent.call_method1(method, args)?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Additional context for log formatting, tracking which request a unit of work
|
||||
/// belongs to and accounting CPU/DB usage against it. Contexts are scoped within
|
||||
/// a `with` block.
|
||||
///
|
||||
/// A native port of the former Python `LoggingContext`; the attribute surface,
|
||||
/// methods, error-message wording and abuse-detection behaviour are preserved so
|
||||
/// callers (and Python subclasses) are unaffected.
|
||||
///
|
||||
/// Construction is deliberately split between `__new__` (which allocates a blank
|
||||
/// instance) and `__init__` (which does the real initialisation), mirroring how a
|
||||
/// pure-Python class behaves. This lets Python subclasses — in particular
|
||||
/// `synapse.metrics.background_process_metrics.BackgroundProcessLoggingContext`,
|
||||
/// which composes a name and then calls `super().__init__(name=..., ...)` — work
|
||||
/// unchanged.
|
||||
#[pyclass(subclass, name = "LoggingContext", module = "synapse.logging.context")]
|
||||
pub struct LoggingContext {
|
||||
/// Name for the context, used in logging.
|
||||
#[pyo3(get, set)]
|
||||
name: String,
|
||||
/// The homeserver name this context is associated with.
|
||||
#[pyo3(get, set)]
|
||||
server_name: String,
|
||||
/// The OS thread id (`threading.get_ident()`) this context was created on;
|
||||
/// activity on any other thread is an error.
|
||||
#[pyo3(get, set)]
|
||||
main_thread: u64,
|
||||
/// Whether `__exit__` has run. Re-activating a finished context is an error.
|
||||
#[pyo3(get, set)]
|
||||
finished: bool,
|
||||
/// The thread resource usage (`resource.struct_rusage`) captured when this
|
||||
/// context became active, or `None` if it is not currently active.
|
||||
#[pyo3(get, set)]
|
||||
usage_start: Option<Py<PyAny>>,
|
||||
/// A short human-readable tag (e.g. the sync type); always a `str`.
|
||||
#[pyo3(get, set)]
|
||||
tag: String,
|
||||
/// The resources used by this context so far. Exposed to Python as
|
||||
/// `_resource_usage` (see the getter below); mutated in place.
|
||||
resource_usage: Py<ContextResourceUsage>,
|
||||
/// The context that was current when this one was created; restored on exit.
|
||||
#[pyo3(get, set)]
|
||||
previous_context: Option<Py<PyAny>>,
|
||||
/// The parent context, if any; usage is propagated up to it.
|
||||
#[pyo3(get, set)]
|
||||
parent_context: Option<Py<PyAny>>,
|
||||
/// The `ContextRequest` this work belongs to, if any.
|
||||
#[pyo3(get, set)]
|
||||
request: Option<Py<PyAny>>,
|
||||
/// The opentracing scope associated with this context, if any.
|
||||
#[pyo3(get, set)]
|
||||
scope: Option<Py<PyAny>>,
|
||||
}
|
||||
|
||||
#[pymethods]
|
||||
impl LoggingContext {
|
||||
/// Allocate a blank context. The real initialisation happens in `__init__`;
|
||||
/// see the type docstring for why this is split. Extra positional/keyword
|
||||
/// arguments are accepted and ignored so that subclasses passing their own
|
||||
/// constructor arguments up through `type.__call__` (which feeds the same
|
||||
/// arguments to both `__new__` and `__init__`) are not rejected here.
|
||||
#[new]
|
||||
#[pyo3(signature = (*_args, **_kwargs))]
|
||||
fn __new__(
|
||||
py: Python<'_>,
|
||||
_args: &Bound<'_, PyTuple>,
|
||||
_kwargs: Option<&Bound<'_, PyDict>>,
|
||||
) -> PyResult<Self> {
|
||||
Ok(LoggingContext {
|
||||
name: String::new(),
|
||||
server_name: String::new(),
|
||||
main_thread: 0,
|
||||
finished: false,
|
||||
usage_start: None,
|
||||
tag: String::new(),
|
||||
resource_usage: Py::new(py, ContextResourceUsage::default())?,
|
||||
previous_context: None,
|
||||
parent_context: None,
|
||||
request: None,
|
||||
scope: None,
|
||||
})
|
||||
}
|
||||
|
||||
#[pyo3(signature = (*, name, server_name, parent_context=None, request=None))]
|
||||
fn __init__(
|
||||
&mut self,
|
||||
py: Python<'_>,
|
||||
name: String,
|
||||
server_name: String,
|
||||
parent_context: Option<Py<PyAny>>,
|
||||
request: Option<Py<PyAny>>,
|
||||
) -> PyResult<()> {
|
||||
self.previous_context = Some(current_context(py));
|
||||
|
||||
// track the resources used by this context so far
|
||||
self.resource_usage = Py::new(py, ContextResourceUsage::default())?;
|
||||
|
||||
// The thread resource usage when the logcontext became active. None if
|
||||
// the context is not currently active.
|
||||
self.usage_start = None;
|
||||
|
||||
self.name = name;
|
||||
self.server_name = server_name;
|
||||
self.main_thread = get_thread_id(py)?;
|
||||
self.request = None;
|
||||
self.tag = String::new();
|
||||
self.scope = None;
|
||||
|
||||
// keep track of whether we have hit the __exit__ block for this context
|
||||
self.finished = false;
|
||||
|
||||
// Inherit some fields from the parent context (read before we move it
|
||||
// into `self`, so no borrow of `self.parent_context` is held).
|
||||
if let Some(parent) = &parent_context {
|
||||
let parent = parent.bind(py);
|
||||
// which request this corresponds to
|
||||
self.request = none_to_option(parent.getattr("request")?);
|
||||
// we also track the current scope
|
||||
self.scope = none_to_option(parent.getattr("scope")?);
|
||||
}
|
||||
|
||||
if let Some(request) = request {
|
||||
// the request param overrides the request from the parent context
|
||||
self.request = Some(request);
|
||||
}
|
||||
|
||||
self.parent_context = parent_context;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// The resources used by this context so far (mutated in place). Named
|
||||
/// `_resource_usage` to match the historical private attribute.
|
||||
#[getter(_resource_usage)]
|
||||
fn get_resource_usage_attr(&self, py: Python<'_>) -> Py<ContextResourceUsage> {
|
||||
self.resource_usage.clone_ref(py)
|
||||
}
|
||||
|
||||
fn __str__(&self) -> String {
|
||||
self.name.clone()
|
||||
}
|
||||
|
||||
/// Enter this logging context, making it the current context.
|
||||
fn __enter__<'py>(slf: Bound<'py, Self>) -> PyResult<Bound<'py, Self>> {
|
||||
let py = slf.py();
|
||||
let module = context_module(py)?;
|
||||
|
||||
let (name, previous) = {
|
||||
let this = slf.borrow();
|
||||
(
|
||||
this.name.clone(),
|
||||
this.previous_context.as_ref().map(|p| p.clone_ref(py)),
|
||||
)
|
||||
};
|
||||
|
||||
debug!(target: DEBUG_LOGGER_NAME, "LoggingContext({name}).__enter__");
|
||||
|
||||
let old_context = module.getattr("set_current_context")?.call1((&slf,))?;
|
||||
|
||||
let previous = previous.unwrap_or_else(|| py.None());
|
||||
if previous.bind(py).ne(&old_context)? {
|
||||
let previous_repr: String = previous.bind(py).repr()?.extract()?;
|
||||
let old_repr: String = old_context.repr()?.extract()?;
|
||||
logcontext_error(
|
||||
&module,
|
||||
format!("Expected previous context {previous_repr}, found {old_repr}"),
|
||||
)?;
|
||||
}
|
||||
|
||||
Ok(slf)
|
||||
}
|
||||
|
||||
/// Restore the previous logging context. Returns `None` (does not suppress
|
||||
/// exceptions).
|
||||
fn __exit__(
|
||||
slf: Bound<'_, Self>,
|
||||
_exc_type: Bound<'_, PyAny>,
|
||||
_exc_value: Bound<'_, PyAny>,
|
||||
_traceback: Bound<'_, PyAny>,
|
||||
) -> PyResult<()> {
|
||||
let py = slf.py();
|
||||
let module = context_module(py)?;
|
||||
|
||||
let (name, previous) = {
|
||||
let this = slf.borrow();
|
||||
(
|
||||
this.name.clone(),
|
||||
this.previous_context.as_ref().map(|p| p.clone_ref(py)),
|
||||
)
|
||||
};
|
||||
let previous = previous.unwrap_or_else(|| py.None());
|
||||
|
||||
if log_enabled!(target: DEBUG_LOGGER_NAME, Level::Debug) {
|
||||
// Match the Python `%s`: the str() of the previous context. Computed
|
||||
// only when the opt-in debug logger is actually enabled.
|
||||
let previous_str: String = previous.bind(py).str()?.extract()?;
|
||||
debug!(
|
||||
target: DEBUG_LOGGER_NAME,
|
||||
"LoggingContext({name}).__exit__ --> {previous_str}"
|
||||
);
|
||||
}
|
||||
|
||||
let current = module
|
||||
.getattr("set_current_context")?
|
||||
.call1((previous.bind(py),))?;
|
||||
|
||||
if !current.is(&slf) {
|
||||
if current.is(sentinel(py).bind(py)) {
|
||||
logcontext_error(&module, format!("Expected logging context {name} was lost"))?;
|
||||
} else {
|
||||
let current_str: String = current.str()?.extract()?;
|
||||
logcontext_error(
|
||||
&module,
|
||||
format!("Expected logging context {name} but found {current_str}"),
|
||||
)?;
|
||||
}
|
||||
}
|
||||
|
||||
// the fact that we are here suggests that the caller thinks everything is
|
||||
// done and dusted for this logcontext, and further activity will not get
|
||||
// recorded against the correct metrics.
|
||||
slf.borrow_mut().finished = true;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Record that this logcontext is currently running.
|
||||
///
|
||||
/// This should not be called directly: use `set_current_context`.
|
||||
fn start(slf: Bound<'_, Self>, rusage: Option<Py<PyAny>>) -> PyResult<()> {
|
||||
let py = slf.py();
|
||||
let module = context_module(py)?;
|
||||
let name = slf.borrow().name.clone();
|
||||
let main_thread = slf.borrow().main_thread;
|
||||
|
||||
if get_thread_id(py)? != main_thread {
|
||||
logcontext_error(
|
||||
&module,
|
||||
format!("Started logcontext {name} on different thread"),
|
||||
)?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if slf.borrow().finished {
|
||||
logcontext_error(&module, format!("Re-starting finished log context {name}"))?;
|
||||
}
|
||||
|
||||
// If we haven't already started, record the thread resource usage so far.
|
||||
if is_truthy(py, &slf.borrow().usage_start)? {
|
||||
logcontext_error(
|
||||
&module,
|
||||
format!("Re-starting already-active log context {name}"),
|
||||
)?;
|
||||
} else {
|
||||
slf.borrow_mut().usage_start = rusage;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Record that this logcontext is no longer running.
|
||||
///
|
||||
/// This should not be called directly: use `set_current_context`.
|
||||
fn stop(slf: Bound<'_, Self>, rusage: Option<Py<PyAny>>) -> PyResult<()> {
|
||||
let py = slf.py();
|
||||
let module = context_module(py)?;
|
||||
let name = slf.borrow().name.clone();
|
||||
let main_thread = slf.borrow().main_thread;
|
||||
|
||||
// Mirror Python's `try: ... finally: self.usage_start = None`.
|
||||
let result = (|| -> PyResult<()> {
|
||||
if get_thread_id(py)? != main_thread {
|
||||
logcontext_error(
|
||||
&module,
|
||||
format!("Stopped logcontext {name} on different thread"),
|
||||
)?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if !is_truthy(py, &rusage)? {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Record the cpu used since we started.
|
||||
if !is_truthy(py, &slf.borrow().usage_start)? {
|
||||
logcontext_error(
|
||||
&module,
|
||||
format!("Called stop on logcontext {name} without recording a start rusage"),
|
||||
)?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let current = rusage.as_ref().expect("rusage is truthy").bind(py);
|
||||
let (utime_delta, stime_delta): (f64, f64) =
|
||||
slf.call_method1("_get_cputime", (current,))?.extract()?;
|
||||
slf.call_method1("add_cputime", (utime_delta, stime_delta))?;
|
||||
Ok(())
|
||||
})();
|
||||
|
||||
slf.borrow_mut().usage_start = None;
|
||||
result
|
||||
}
|
||||
|
||||
/// Get a *copy* of the resources used by this logcontext so far.
|
||||
fn get_resource_usage(slf: Bound<'_, Self>) -> PyResult<ContextResourceUsage> {
|
||||
let py = slf.py();
|
||||
|
||||
// we always return a copy, for consistency
|
||||
let mut res = slf.borrow().resource_usage.borrow(py).clone();
|
||||
|
||||
let active = is_truthy(py, &slf.borrow().usage_start)?;
|
||||
|
||||
// If we are on the correct thread and we're currently running then we can
|
||||
// include resource usage so far.
|
||||
let is_main_thread = get_thread_id(py)? == slf.borrow().main_thread;
|
||||
if active && is_main_thread {
|
||||
let rusage = context_module(py)?
|
||||
.getattr("get_thread_resource_usage")?
|
||||
.call0()?;
|
||||
if rusage.is_none() {
|
||||
return Err(PyAssertionError::new_err(
|
||||
"get_thread_resource_usage() returned None while active",
|
||||
));
|
||||
}
|
||||
let (utime_delta, stime_delta): (f64, f64) =
|
||||
slf.call_method1("_get_cputime", (rusage,))?.extract()?;
|
||||
res.ru_utime += utime_delta;
|
||||
res.ru_stime += stime_delta;
|
||||
}
|
||||
|
||||
Ok(res)
|
||||
}
|
||||
|
||||
/// Get the cpu usage time between `start()` and the given rusage. Returns
|
||||
/// `(seconds in user mode, seconds in system mode)`.
|
||||
fn _get_cputime(&self, py: Python<'_>, current: Bound<'_, PyAny>) -> PyResult<(f64, f64)> {
|
||||
let usage_start = self
|
||||
.usage_start
|
||||
.as_ref()
|
||||
.ok_or_else(|| PyAssertionError::new_err("usage_start is None"))?;
|
||||
let usage_start = usage_start.bind(py);
|
||||
|
||||
let cur_utime: f64 = current.getattr("ru_utime")?.extract()?;
|
||||
let cur_stime: f64 = current.getattr("ru_stime")?.extract()?;
|
||||
let start_utime: f64 = usage_start.getattr("ru_utime")?.extract()?;
|
||||
let start_stime: f64 = usage_start.getattr("ru_stime")?.extract()?;
|
||||
|
||||
let mut utime_delta = cur_utime - start_utime;
|
||||
let mut stime_delta = cur_stime - start_stime;
|
||||
|
||||
// sanity check
|
||||
if utime_delta < 0.0 {
|
||||
context_module(py)?.getattr("logger")?.call_method1(
|
||||
"error",
|
||||
("utime went backwards! %f < %f", cur_utime, start_utime),
|
||||
)?;
|
||||
utime_delta = 0.0;
|
||||
}
|
||||
if stime_delta < 0.0 {
|
||||
context_module(py)?.getattr("logger")?.call_method1(
|
||||
"error",
|
||||
("stime went backwards! %f < %f", cur_stime, start_stime),
|
||||
)?;
|
||||
stime_delta = 0.0;
|
||||
}
|
||||
|
||||
Ok((utime_delta, stime_delta))
|
||||
}
|
||||
|
||||
/// Update the CPU time usage of this context (and any parents, recursively).
|
||||
fn add_cputime(&self, py: Python<'_>, utime_delta: f64, stime_delta: f64) -> PyResult<()> {
|
||||
{
|
||||
let mut usage = self.resource_usage.borrow_mut(py);
|
||||
usage.ru_utime += utime_delta;
|
||||
usage.ru_stime += stime_delta;
|
||||
}
|
||||
forward_to_parent(
|
||||
&self.parent_context,
|
||||
py,
|
||||
"add_cputime",
|
||||
(utime_delta, stime_delta),
|
||||
)
|
||||
}
|
||||
|
||||
/// Record the use of a database transaction and how long it took.
|
||||
fn add_database_transaction(&self, py: Python<'_>, duration_sec: f64) -> PyResult<()> {
|
||||
if duration_sec < 0.0 {
|
||||
return Err(PyValueError::new_err(
|
||||
"DB txn time can only be non-negative",
|
||||
));
|
||||
}
|
||||
{
|
||||
let mut usage = self.resource_usage.borrow_mut(py);
|
||||
usage.db_txn_count += 1;
|
||||
usage.db_txn_duration_sec += duration_sec;
|
||||
}
|
||||
forward_to_parent(
|
||||
&self.parent_context,
|
||||
py,
|
||||
"add_database_transaction",
|
||||
(duration_sec,),
|
||||
)
|
||||
}
|
||||
|
||||
/// Record a use of the database pool (the time taken to get a connection).
|
||||
fn add_database_scheduled(&self, py: Python<'_>, sched_sec: f64) -> PyResult<()> {
|
||||
if sched_sec < 0.0 {
|
||||
return Err(PyValueError::new_err(
|
||||
"DB scheduling time can only be non-negative",
|
||||
));
|
||||
}
|
||||
{
|
||||
let mut usage = self.resource_usage.borrow_mut(py);
|
||||
usage.db_sched_duration_sec += sched_sec;
|
||||
}
|
||||
forward_to_parent(
|
||||
&self.parent_context,
|
||||
py,
|
||||
"add_database_scheduled",
|
||||
(sched_sec,),
|
||||
)
|
||||
}
|
||||
|
||||
/// Record a number of events being fetched from the db.
|
||||
fn record_event_fetch(&self, py: Python<'_>, event_count: i64) -> PyResult<()> {
|
||||
{
|
||||
let mut usage = self.resource_usage.borrow_mut(py);
|
||||
usage.evt_db_fetch_count += event_count;
|
||||
}
|
||||
forward_to_parent(
|
||||
&self.parent_context,
|
||||
py,
|
||||
"record_event_fetch",
|
||||
(event_count,),
|
||||
)
|
||||
}
|
||||
|
||||
/// Traverse referenced Python objects for the cyclic garbage collector.
|
||||
/// `scope` and the context can reference each other, forming a real cycle.
|
||||
fn __traverse__(&self, visit: PyVisit<'_>) -> Result<(), PyTraverseError> {
|
||||
if let Some(previous_context) = &self.previous_context {
|
||||
visit.call(previous_context)?;
|
||||
}
|
||||
if let Some(parent_context) = &self.parent_context {
|
||||
visit.call(parent_context)?;
|
||||
}
|
||||
if let Some(request) = &self.request {
|
||||
visit.call(request)?;
|
||||
}
|
||||
if let Some(scope) = &self.scope {
|
||||
visit.call(scope)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn __clear__(&mut self) {
|
||||
self.previous_context = None;
|
||||
self.parent_context = None;
|
||||
self.request = None;
|
||||
self.scope = None;
|
||||
}
|
||||
}
|
||||
|
||||
/// Register the Python sentinel logcontext.
|
||||
///
|
||||
/// Called once from `synapse.logging.context` at import time. Registering twice
|
||||
@@ -299,9 +848,11 @@ pub fn swap_current_context(py: Python<'_>, context: Py<PyAny>) -> Py<PyAny> {
|
||||
pub fn register_module(py: Python<'_>, m: &Bound<'_, PyModule>) -> PyResult<()> {
|
||||
let child_module: Bound<'_, PyModule> = PyModule::new(py, "logcontext")?;
|
||||
child_module.add_class::<ContextResourceUsage>()?;
|
||||
child_module.add_class::<LoggingContext>()?;
|
||||
child_module.add_function(wrap_pyfunction!(current_context, &child_module)?)?;
|
||||
child_module.add_function(wrap_pyfunction!(swap_current_context, &child_module)?)?;
|
||||
child_module.add_function(wrap_pyfunction!(register_sentinel, &child_module)?)?;
|
||||
child_module.add("DEBUG_LOGGER_NAME", DEBUG_LOGGER_NAME)?;
|
||||
|
||||
m.add_submodule(&child_module)?;
|
||||
|
||||
|
||||
+3
-292
@@ -31,7 +31,6 @@ See doc/log_contexts.rst for details on how this works.
|
||||
"""
|
||||
|
||||
import logging
|
||||
import threading
|
||||
import typing
|
||||
from types import TracebackType
|
||||
from typing import (
|
||||
@@ -53,7 +52,9 @@ from twisted.python.threadpool import ThreadPool
|
||||
|
||||
from synapse.logging.loggers import ExplicitlyConfiguredLogger
|
||||
from synapse.synapse_rust.logcontext import (
|
||||
DEBUG_LOGGER_NAME,
|
||||
ContextResourceUsage as ContextResourceUsage,
|
||||
LoggingContext as LoggingContext,
|
||||
current_context as current_context,
|
||||
register_sentinel,
|
||||
swap_current_context,
|
||||
@@ -61,14 +62,13 @@ from synapse.synapse_rust.logcontext import (
|
||||
from synapse.util.stringutils import random_string_insecure_fast
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from synapse.logging.scopecontextmanager import _LogContextScope
|
||||
from synapse.types import ISynapseReactor
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
original_logger_class = logging.getLoggerClass()
|
||||
logging.setLoggerClass(ExplicitlyConfiguredLogger)
|
||||
logcontext_debug_logger = logging.getLogger("synapse.logging.context.debug")
|
||||
logcontext_debug_logger = logging.getLogger(DEBUG_LOGGER_NAME)
|
||||
"""
|
||||
A logger for debugging when the logcontext switches.
|
||||
|
||||
@@ -110,15 +110,6 @@ def logcontext_error(msg: str) -> None:
|
||||
logger.warning(msg)
|
||||
|
||||
|
||||
# get an id for the current thread.
|
||||
#
|
||||
# threading.get_ident doesn't actually return an OS-level tid, and annoyingly,
|
||||
# on Linux it actually returns the same value either side of a fork() call. However
|
||||
# we only fork in one place, so it's not worth the hoop-jumping to get a real tid.
|
||||
#
|
||||
get_thread_id = threading.get_ident
|
||||
|
||||
|
||||
@attr.s(slots=True, auto_attribs=True)
|
||||
class ContextRequest:
|
||||
"""
|
||||
@@ -207,286 +198,6 @@ SENTINEL_CONTEXT = _Sentinel()
|
||||
register_sentinel(SENTINEL_CONTEXT)
|
||||
|
||||
|
||||
class LoggingContext:
|
||||
"""Additional context for log formatting. Contexts are scoped within a
|
||||
"with" block.
|
||||
|
||||
If a parent is given when creating a new context, then:
|
||||
- logging fields are copied from the parent to the new context on entry
|
||||
- when the new context exits, the cpu usage stats are copied from the
|
||||
child to the parent
|
||||
|
||||
Args:
|
||||
name: Name for the context for logging.
|
||||
server_name: The name of the server this context is associated with
|
||||
(`config.server.server_name` or `hs.hostname`)
|
||||
parent_context (LoggingContext|None): The parent of the new context
|
||||
request: Synapse Request Context object. Useful to associate all the logs
|
||||
happening to a given request.
|
||||
|
||||
"""
|
||||
|
||||
__slots__ = [
|
||||
"previous_context",
|
||||
"name",
|
||||
"server_name",
|
||||
"parent_context",
|
||||
"_resource_usage",
|
||||
"usage_start",
|
||||
"main_thread",
|
||||
"finished",
|
||||
"request",
|
||||
"tag",
|
||||
"scope",
|
||||
]
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
name: str,
|
||||
server_name: str,
|
||||
parent_context: "LoggingContext | None" = None,
|
||||
request: ContextRequest | None = None,
|
||||
) -> None:
|
||||
self.previous_context = current_context()
|
||||
|
||||
# track the resources used by this context so far
|
||||
self._resource_usage = ContextResourceUsage()
|
||||
|
||||
# The thread resource usage when the logcontext became active. None
|
||||
# if the context is not currently active.
|
||||
self.usage_start: resource.struct_rusage | None = None
|
||||
|
||||
self.name = name
|
||||
self.server_name = server_name
|
||||
self.main_thread = get_thread_id()
|
||||
self.request = None
|
||||
self.tag = ""
|
||||
self.scope: "_LogContextScope" | None = None
|
||||
|
||||
# keep track of whether we have hit the __exit__ block for this context
|
||||
# (suggesting that the the thing that created the context thinks it should
|
||||
# be finished, and that re-activating it would suggest an error).
|
||||
self.finished = False
|
||||
|
||||
self.parent_context = parent_context
|
||||
|
||||
# Inherit some fields from the parent context
|
||||
if self.parent_context is not None:
|
||||
# which request this corresponds to
|
||||
self.request = self.parent_context.request
|
||||
|
||||
# we also track the current scope:
|
||||
self.scope = self.parent_context.scope
|
||||
|
||||
if request is not None:
|
||||
# the request param overrides the request from the parent context
|
||||
self.request = request
|
||||
|
||||
def __str__(self) -> str:
|
||||
return self.name
|
||||
|
||||
def __enter__(self) -> "LoggingContext":
|
||||
"""Enters this logging context into thread local storage"""
|
||||
logcontext_debug_logger.debug("LoggingContext(%s).__enter__", self.name)
|
||||
old_context = set_current_context(self)
|
||||
if self.previous_context != old_context:
|
||||
logcontext_error(
|
||||
"Expected previous context %r, found %r"
|
||||
% (
|
||||
self.previous_context,
|
||||
old_context,
|
||||
)
|
||||
)
|
||||
return self
|
||||
|
||||
def __exit__(
|
||||
self,
|
||||
type: type[BaseException] | None,
|
||||
value: BaseException | None,
|
||||
traceback: TracebackType | None,
|
||||
) -> None:
|
||||
"""Restore the logging context in thread local storage to the state it
|
||||
was before this context was entered.
|
||||
Returns:
|
||||
None to avoid suppressing any exceptions that were thrown.
|
||||
"""
|
||||
logcontext_debug_logger.debug(
|
||||
"LoggingContext(%s).__exit__ --> %s", self.name, self.previous_context
|
||||
)
|
||||
current = set_current_context(self.previous_context)
|
||||
if current is not self:
|
||||
if current is SENTINEL_CONTEXT:
|
||||
logcontext_error("Expected logging context %s was lost" % (self,))
|
||||
else:
|
||||
logcontext_error(
|
||||
"Expected logging context %s but found %s" % (self, current)
|
||||
)
|
||||
|
||||
# the fact that we are here suggests that the caller thinks that everything
|
||||
# is done and dusted for this logcontext, and further activity will not get
|
||||
# recorded against the correct metrics.
|
||||
self.finished = True
|
||||
|
||||
def start(self, rusage: "resource.struct_rusage | None") -> None:
|
||||
"""
|
||||
Record that this logcontext is currently running.
|
||||
|
||||
This should not be called directly: use set_current_context
|
||||
|
||||
Args:
|
||||
rusage: the resources used by the current thread, at the point of
|
||||
switching to this logcontext. May be None if this platform doesn't
|
||||
support getrusuage.
|
||||
"""
|
||||
if get_thread_id() != self.main_thread:
|
||||
logcontext_error("Started logcontext %s on different thread" % (self,))
|
||||
return
|
||||
|
||||
if self.finished:
|
||||
logcontext_error("Re-starting finished log context %s" % (self,))
|
||||
|
||||
# If we haven't already started record the thread resource usage so
|
||||
# far
|
||||
if self.usage_start:
|
||||
logcontext_error("Re-starting already-active log context %s" % (self,))
|
||||
else:
|
||||
self.usage_start = rusage
|
||||
|
||||
def stop(self, rusage: "resource.struct_rusage | None") -> None:
|
||||
"""
|
||||
Record that this logcontext is no longer running.
|
||||
|
||||
This should not be called directly: use set_current_context
|
||||
|
||||
Args:
|
||||
rusage: the resources used by the current thread, at the point of
|
||||
switching away from this logcontext. May be None if this platform
|
||||
doesn't support getrusuage.
|
||||
"""
|
||||
|
||||
try:
|
||||
if get_thread_id() != self.main_thread:
|
||||
logcontext_error("Stopped logcontext %s on different thread" % (self,))
|
||||
return
|
||||
|
||||
if not rusage:
|
||||
return
|
||||
|
||||
# Record the cpu used since we started
|
||||
if not self.usage_start:
|
||||
logcontext_error(
|
||||
"Called stop on logcontext %s without recording a start rusage"
|
||||
% (self,)
|
||||
)
|
||||
return
|
||||
|
||||
utime_delta, stime_delta = self._get_cputime(rusage)
|
||||
self.add_cputime(utime_delta, stime_delta)
|
||||
finally:
|
||||
self.usage_start = None
|
||||
|
||||
def get_resource_usage(self) -> ContextResourceUsage:
|
||||
"""Get resources used by this logcontext so far.
|
||||
|
||||
Returns:
|
||||
A *copy* of the object tracking resource usage so far
|
||||
"""
|
||||
# we always return a copy, for consistency
|
||||
res = self._resource_usage.copy()
|
||||
|
||||
# If we are on the correct thread and we're currently running then we
|
||||
# can include resource usage so far.
|
||||
is_main_thread = get_thread_id() == self.main_thread
|
||||
if self.usage_start and is_main_thread:
|
||||
rusage = get_thread_resource_usage()
|
||||
assert rusage is not None
|
||||
utime_delta, stime_delta = self._get_cputime(rusage)
|
||||
res.ru_utime += utime_delta
|
||||
res.ru_stime += stime_delta
|
||||
|
||||
return res
|
||||
|
||||
def _get_cputime(self, current: "resource.struct_rusage") -> tuple[float, float]:
|
||||
"""Get the cpu usage time between start() and the given rusage
|
||||
|
||||
Args:
|
||||
rusage: the current resource usage
|
||||
|
||||
Returns: tuple[float, float]: seconds in user mode, seconds in system mode
|
||||
"""
|
||||
assert self.usage_start is not None
|
||||
|
||||
utime_delta = current.ru_utime - self.usage_start.ru_utime
|
||||
stime_delta = current.ru_stime - self.usage_start.ru_stime
|
||||
|
||||
# sanity check
|
||||
if utime_delta < 0:
|
||||
logger.error(
|
||||
"utime went backwards! %f < %f",
|
||||
current.ru_utime,
|
||||
self.usage_start.ru_utime,
|
||||
)
|
||||
utime_delta = 0
|
||||
|
||||
if stime_delta < 0:
|
||||
logger.error(
|
||||
"stime went backwards! %f < %f",
|
||||
current.ru_stime,
|
||||
self.usage_start.ru_stime,
|
||||
)
|
||||
stime_delta = 0
|
||||
|
||||
return utime_delta, stime_delta
|
||||
|
||||
def add_cputime(self, utime_delta: float, stime_delta: float) -> None:
|
||||
"""Update the CPU time usage of this context (and any parents, recursively).
|
||||
|
||||
Args:
|
||||
utime_delta: additional user time, in seconds, spent in this context.
|
||||
stime_delta: additional system time, in seconds, spent in this context.
|
||||
"""
|
||||
self._resource_usage.ru_utime += utime_delta
|
||||
self._resource_usage.ru_stime += stime_delta
|
||||
if self.parent_context:
|
||||
self.parent_context.add_cputime(utime_delta, stime_delta)
|
||||
|
||||
def add_database_transaction(self, duration_sec: float) -> None:
|
||||
"""Record the use of a database transaction and the length of time it took.
|
||||
|
||||
Args:
|
||||
duration_sec: The number of seconds the database transaction took.
|
||||
"""
|
||||
if duration_sec < 0:
|
||||
raise ValueError("DB txn time can only be non-negative")
|
||||
self._resource_usage.db_txn_count += 1
|
||||
self._resource_usage.db_txn_duration_sec += duration_sec
|
||||
if self.parent_context:
|
||||
self.parent_context.add_database_transaction(duration_sec)
|
||||
|
||||
def add_database_scheduled(self, sched_sec: float) -> None:
|
||||
"""Record a use of the database pool
|
||||
|
||||
Args:
|
||||
sched_sec: number of seconds it took us to get a connection
|
||||
"""
|
||||
if sched_sec < 0:
|
||||
raise ValueError("DB scheduling time can only be non-negative")
|
||||
self._resource_usage.db_sched_duration_sec += sched_sec
|
||||
if self.parent_context:
|
||||
self.parent_context.add_database_scheduled(sched_sec)
|
||||
|
||||
def record_event_fetch(self, event_count: int) -> None:
|
||||
"""Record a number of events being fetched from the db
|
||||
|
||||
Args:
|
||||
event_count: number of events being fetched
|
||||
"""
|
||||
self._resource_usage.evt_db_fetch_count += event_count
|
||||
if self.parent_context:
|
||||
self.parent_context.record_event_fetch(event_count)
|
||||
|
||||
|
||||
class LoggingContextFilter(logging.Filter):
|
||||
"""Logging filter that adds values from the current logging context to each
|
||||
record.
|
||||
|
||||
@@ -10,9 +10,19 @@
|
||||
# See the GNU Affero General Public License for more details:
|
||||
# <https://www.gnu.org/licenses/agpl-3.0.html>.
|
||||
|
||||
from typing import Optional
|
||||
import resource
|
||||
from types import TracebackType
|
||||
from typing import TYPE_CHECKING, Optional
|
||||
|
||||
from synapse.logging.context import LoggingContextOrSentinel
|
||||
from synapse.logging.context import ContextRequest, LoggingContextOrSentinel
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from synapse.logging.scopecontextmanager import _LogContextScope
|
||||
|
||||
DEBUG_LOGGER_NAME: str
|
||||
"""Name of the opt-in logger for logcontext switch tracing
|
||||
(`synapse.logging.context.debug`). Shared with the Rust `debug!` target so the
|
||||
names cannot drift."""
|
||||
|
||||
class ContextResourceUsage:
|
||||
"""Tracks the resources used by a log context."""
|
||||
@@ -32,6 +42,55 @@ class ContextResourceUsage:
|
||||
def __add__(self, other: "ContextResourceUsage") -> "ContextResourceUsage": ...
|
||||
def __sub__(self, other: "ContextResourceUsage") -> "ContextResourceUsage": ...
|
||||
|
||||
class LoggingContext:
|
||||
"""Additional context for log formatting. Contexts are scoped within a
|
||||
"with" block.
|
||||
|
||||
If a parent is given when creating a new context, then:
|
||||
- logging fields are copied from the parent to the new context on entry
|
||||
- when the new context exits, the cpu usage stats are copied from the
|
||||
child to the parent
|
||||
"""
|
||||
|
||||
previous_context: LoggingContextOrSentinel
|
||||
name: str
|
||||
server_name: str
|
||||
parent_context: "Optional[LoggingContext]"
|
||||
usage_start: "Optional[resource.struct_rusage]"
|
||||
main_thread: int
|
||||
finished: bool
|
||||
request: Optional[ContextRequest]
|
||||
tag: str
|
||||
scope: "Optional[_LogContextScope]"
|
||||
_resource_usage: ContextResourceUsage
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
name: str,
|
||||
server_name: str,
|
||||
parent_context: "Optional[LoggingContext]" = None,
|
||||
request: Optional[ContextRequest] = None,
|
||||
) -> None: ...
|
||||
def __str__(self) -> str: ...
|
||||
def __enter__(self) -> "LoggingContext": ...
|
||||
def __exit__(
|
||||
self,
|
||||
type: Optional[type[BaseException]],
|
||||
value: Optional[BaseException],
|
||||
traceback: Optional[TracebackType],
|
||||
) -> None: ...
|
||||
def start(self, rusage: "Optional[resource.struct_rusage]") -> None: ...
|
||||
def stop(self, rusage: "Optional[resource.struct_rusage]") -> None: ...
|
||||
def get_resource_usage(self) -> ContextResourceUsage: ...
|
||||
def _get_cputime(
|
||||
self, current: "resource.struct_rusage"
|
||||
) -> tuple[float, float]: ...
|
||||
def add_cputime(self, utime_delta: float, stime_delta: float) -> None: ...
|
||||
def add_database_transaction(self, duration_sec: float) -> None: ...
|
||||
def add_database_scheduled(self, sched_sec: float) -> None: ...
|
||||
def record_event_fetch(self, event_count: int) -> None: ...
|
||||
|
||||
def current_context() -> LoggingContextOrSentinel:
|
||||
"""Get the current logging context.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user