diff --git a/rust/src/logging/context.rs b/rust/src/logging/context.rs index 3afe3ccf8f..59c62d4df8 100644 --- a/rust/src/logging/context.rs +++ b/rust/src/logging/context.rs @@ -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> = 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> { + 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> = OnceCell::new(); + +/// The current OS thread id, matching Python's `threading.get_ident()`. +fn get_thread_id(py: Python<'_>) -> PyResult { + let get_ident = GET_THREAD_ID.get_or_try_init(|| -> PyResult> { + 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> { + 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>) -> PyResult { + 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: 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>, + /// 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, + /// The context that was current when this one was created; restored on exit. + #[pyo3(get, set)] + previous_context: Option>, + /// The parent context, if any; usage is propagated up to it. + #[pyo3(get, set)] + parent_context: Option>, + /// The `ContextRequest` this work belongs to, if any. + #[pyo3(get, set)] + request: Option>, + /// The opentracing scope associated with this context, if any. + #[pyo3(get, set)] + scope: Option>, +} + +#[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 { + 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>, + request: Option>, + ) -> 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 { + 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> { + 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>) -> 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>) -> 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 { + 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) -> Py { pub fn register_module(py: Python<'_>, m: &Bound<'_, PyModule>) -> PyResult<()> { let child_module: Bound<'_, PyModule> = PyModule::new(py, "logcontext")?; child_module.add_class::()?; + child_module.add_class::()?; 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)?; diff --git a/synapse/logging/context.py b/synapse/logging/context.py index e25e9355b6..255802eba1 100644 --- a/synapse/logging/context.py +++ b/synapse/logging/context.py @@ -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. diff --git a/synapse/synapse_rust/logcontext.pyi b/synapse/synapse_rust/logcontext.pyi index c0c79c16a7..7a89bb8c28 100644 --- a/synapse/synapse_rust/logcontext.pyi +++ b/synapse/synapse_rust/logcontext.pyi @@ -10,9 +10,19 @@ # See the GNU Affero General Public License for more details: # . -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.