mirror of
https://github.com/element-hq/synapse.git
synced 2026-08-14 15:50:19 +00:00
Move the logcontext storage and LoggingContext to Rust
The "current logcontext" slot moves from a Python threading.local into the Rust extension, together with a native port of LoggingContext. A Python thread-local is invisible to Rust: each tokio worker thread would see its own slot, permanently at the sentinel, so logging emitted from Rust could not be attributed to the request that caused it. This lays the storage groundwork; a follow-up change gives tokio tasks a task-scoped capture (see the module-doc TODO). Design notes, for review: - The slot is typed: Option<Py<LoggingContext>>, with None representing the sentinel. The _Sentinel class and SENTINEL_CONTEXT singleton stay pure Python, unchanged; thin wrappers on current_context and set_current_context in synapse.logging.context convert between the singleton and None at the boundary, so no Rust code ever sees or produces the sentinel object. pyo3's extraction enforces the type: anything that is not a LoggingContext (or subclass) or None raises TypeError. - The accounting policy is native too: set_current_context reads the thread rusage once via libc (no per-switch struct_rusage allocation) and runs the stop/start bookkeeping inline for base LoggingContexts, only dispatching through Python for subclasses (BackgroundProcessLoggingContext) so their overrides run. start()/stop() now take an Optional[tuple[float, float]] instead of a struct_rusage, the get_thread_resource_usage/is_thread_resource_usage_supported/ get_thread_id module helpers are gone, LoggingContext.previous_context is now Optional[LoggingContext] (None where it used to hold SENTINEL_CONTEXT), and the nominally-private _resource_usage attribute is no longer exposed (nothing read it; use get_resource_usage()) — worth an upgrade note when this is released, as out-of-tree code may rely on the old shapes. - The switch path avoids per-operation allocation and Python round-trips: names are stored as Py<PyString> (LoggingContextFilter reads server_name and str(context) per log record process-wide, now INCREF-only), error messages materialise the context name only in the cold branches, and the thread id is read via PyThread_get_thread_ident (the exact value threading.get_ident() returns) rather than by calling into Python. - The attribute surface, method set and error-message wording are a compatibility contract, pinned by the characterization tests (which now exercise the tuple-based start/stop API). - The opt-in synapse.logging.context.debug switch traces are emitted from Rust via pyo3-log, whose level cache only refreshes on reset_logging_config(); docs/log_contexts.md documents the manhole procedure. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JFbRtswu7rsHrttJFauUUb
This commit is contained in:
co-authored by
Claude Fable 5
parent
49c163679d
commit
c3dfbfc6ef
Generated
+1
@@ -1342,6 +1342,7 @@ dependencies = [
|
||||
"icu_segmenter",
|
||||
"itertools",
|
||||
"lazy_static",
|
||||
"libc",
|
||||
"log",
|
||||
"mime",
|
||||
"once_cell",
|
||||
|
||||
@@ -564,3 +564,18 @@ loggers:
|
||||
synapse.logging.context.debug:
|
||||
level: DEBUG
|
||||
```
|
||||
|
||||
Note that some of these traces (`LoggingContext(...).__enter__`/`__exit__`) are
|
||||
emitted from Rust via `pyo3-log`, which caches logger levels for performance.
|
||||
Configuring the logger in the logging config (as above) works — the cache is
|
||||
flushed whenever the config is (re)loaded, including on `SIGHUP` — but enabling
|
||||
the logger *at runtime* (e.g. `setLevel(logging.DEBUG)` from the manhole) will
|
||||
not surface the Rust-emitted traces until you also flush the cache:
|
||||
|
||||
```python
|
||||
import logging
|
||||
from synapse.synapse_rust import reset_logging_config
|
||||
|
||||
logging.getLogger("synapse.logging.context.debug").setLevel(logging.DEBUG)
|
||||
reset_logging_config()
|
||||
```
|
||||
|
||||
@@ -29,6 +29,7 @@ bytes = "1.6.0"
|
||||
headers = "0.4.0"
|
||||
http = "1.1.0"
|
||||
lazy_static = "1.4.0"
|
||||
libc = "0.2.174"
|
||||
log = "0.4.17"
|
||||
mime = "0.3.17"
|
||||
pyo3 = { version = "0.28.3", features = [
|
||||
|
||||
+832
-6
@@ -13,16 +13,64 @@
|
||||
*
|
||||
*/
|
||||
|
||||
//! Native counterparts to `synapse.logging.context`.
|
||||
//! Native storage for the Synapse "current logcontext".
|
||||
//!
|
||||
//! Currently just [`ContextResourceUsage`]; the goal is per-operation resource
|
||||
//! accounting with no Python allocation on the switch path.
|
||||
//! The storage lives in Rust rather than in a Python `threading.local` because a
|
||||
//! Python thread-local is invisible to Rust: each tokio worker thread would see
|
||||
//! its own slot, permanently at the sentinel, and logging emitted from Rust —
|
||||
//! including from spawned tokio tasks — could not be attributed to the request
|
||||
//! that caused it.
|
||||
//!
|
||||
//! TODO: the storage for the "current" logcontext, `LoggingContext` itself and
|
||||
//! the sentinel follow — see `synapse.logging.context` for the Python
|
||||
//! implementations being replaced.
|
||||
//! This module holds that storage — a per-OS-thread slot
|
||||
//! ([`THREAD_LOCAL_CONTEXT`]) used by the reactor thread and any
|
||||
//! reactor-managed threadpool threads — along with the logcontext classes
|
||||
//! themselves. `LoggingContextFilter` (and therefore `pyo3-log`) resolves the
|
||||
//! context by calling [`current_context`] at log-record time.
|
||||
//!
|
||||
//! The slot holds an `Option<Py<LoggingContext>>`: `None` means "no context" —
|
||||
//! what Synapse calls the sentinel. The `_Sentinel` marker object itself is pure
|
||||
//! Python (`synapse.logging.context.SENTINEL_CONTEXT`); the wrappers there
|
||||
//! convert between it and `None` at the boundary, so no Rust code ever sees or
|
||||
//! produces the sentinel object.
|
||||
//!
|
||||
//! The accounting policy is native too: [`set_current_context`] reads the thread
|
||||
//! rusage via libc, runs the `stop`/`start` bookkeeping, and uses
|
||||
//! [`swap_current_context`] for the raw slot write.
|
||||
//!
|
||||
//! TODO: tokio tasks do not yet see a logcontext — a worker thread's slot is
|
||||
//! always empty, so Rust-emitted log records still land in the sentinel. A
|
||||
//! task-scoped capture (carried with the task across `.await` points and
|
||||
//! consulted by [`current_context`] ahead of the thread slot) follows in the
|
||||
//! next change.
|
||||
|
||||
use std::cell::RefCell;
|
||||
|
||||
use log::{debug, error, log_enabled, Level};
|
||||
use pyo3::call::PyCallArgs;
|
||||
use pyo3::exceptions::PyValueError;
|
||||
use pyo3::prelude::*;
|
||||
use pyo3::types::{PyDict, PyString, PyTuple};
|
||||
use pyo3::{intern, PyTraverseError, PyVisit};
|
||||
|
||||
/// 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. The slot is typed: it holds a
|
||||
/// [`LoggingContext`] (possibly a Python subclass instance), and `None` means
|
||||
/// "no context set on this thread", i.e. the sentinel — a fresh thread (e.g.
|
||||
/// a new threadpool worker) therefore starts in the sentinel.
|
||||
static THREAD_LOCAL_CONTEXT: RefCell<Option<Py<LoggingContext>>> = const { RefCell::new(None) };
|
||||
}
|
||||
|
||||
/// Tracks the resources used by a log context.
|
||||
///
|
||||
@@ -137,10 +185,716 @@ impl ContextResourceUsage {
|
||||
}
|
||||
}
|
||||
|
||||
/// Call the (possibly test-patched) module-level `logcontext_error(msg)`.
|
||||
fn logcontext_error(py: Python<'_>, msg: String) -> PyResult<()> {
|
||||
let module = py.import("synapse.logging.context")?;
|
||||
module.getattr("logcontext_error")?.call1((msg,))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
extern "C" {
|
||||
/// CPython's thread identifier (`pythread.h`) — the exact value
|
||||
/// `threading.get_ident()` returns. Part of the stable ABI, but not bound
|
||||
/// by pyo3-ffi, so declared here.
|
||||
fn PyThread_get_thread_ident() -> std::os::raw::c_ulong;
|
||||
}
|
||||
|
||||
/// This thread's `threading.get_ident()` value, read natively (`get_ident` is
|
||||
/// a thin wrapper around `PyThread_get_thread_ident`, which is `pthread_self()`
|
||||
/// on POSIX) rather than by calling into Python.
|
||||
///
|
||||
/// Note that `get_ident` is *not* an OS-level tid: on Linux it returns the same
|
||||
/// value either side of a `fork()` call. Synapse forks in exactly one place, so
|
||||
/// contexts created before the fork still pass the `main_thread` affinity check
|
||||
/// after it.
|
||||
fn get_thread_id() -> u64 {
|
||||
// SAFETY: no preconditions; returns an identifier for the calling thread.
|
||||
(unsafe { PyThread_get_thread_ident() }) as u64
|
||||
}
|
||||
|
||||
/// Propagate a usage update to the parent context, if there is a (truthy) one.
|
||||
///
|
||||
/// A plain base `LoggingContext` parent — the overwhelmingly common case, e.g.
|
||||
/// every `Measure`-created nested context — is dispatched via `native` directly,
|
||||
/// skipping the Python method call (attribute lookup, args tuple, call frame)
|
||||
/// that `add_cputime` would otherwise pay on every switch away from a parented
|
||||
/// context and `add_database_*` per DB operation. Anything else truthy is a
|
||||
/// Python subclass, dispatched via `call_method1` so its overrides are
|
||||
/// respected — the same split as `switch_context`. The truthiness guard matches
|
||||
/// Python's `if self.parent_context:` and cannot be relaxed to `is_some()`: the
|
||||
/// sentinel is falsy and implements no `add_cputime`, so that would call a
|
||||
/// nonexistent method on it.
|
||||
fn forward_to_parent<'py, N>(
|
||||
parent: &Option<Py<PyAny>>,
|
||||
py: Python<'py>,
|
||||
method: &str,
|
||||
args: impl PyCallArgs<'py>,
|
||||
native: N,
|
||||
) -> PyResult<()>
|
||||
where
|
||||
N: FnOnce(&Bound<'py, LoggingContext>) -> PyResult<()>,
|
||||
{
|
||||
if let Some(parent) = parent {
|
||||
let parent = parent.bind(py);
|
||||
if let Ok(base) = parent.cast_exact::<LoggingContext>() {
|
||||
return native(base);
|
||||
}
|
||||
if parent.is_truthy()? {
|
||||
parent.call_method1(method, args)?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// The current thread's CPU usage as `(ru_utime, ru_stime)` in seconds, read
|
||||
/// directly via `getrusage(RUSAGE_THREAD)`.
|
||||
///
|
||||
/// Returns `None` where per-thread rusage isn't available — which we take to be
|
||||
/// any non-Linux target (`RUSAGE_THREAD` is Linux-only; macOS gets no per-context
|
||||
/// CPU accounting). Reading it natively avoids allocating a Python
|
||||
/// `resource.struct_rusage` object on every switch.
|
||||
#[cfg(target_os = "linux")]
|
||||
fn get_thread_rusage() -> Option<(f64, f64)> {
|
||||
fn timeval_to_secs(tv: libc::timeval) -> f64 {
|
||||
tv.tv_sec as f64 + tv.tv_usec as f64 / 1_000_000.0
|
||||
}
|
||||
|
||||
// SAFETY: `getrusage` only writes into `usage`, and we only read it once it
|
||||
// reports success.
|
||||
let mut usage = std::mem::MaybeUninit::<libc::rusage>::uninit();
|
||||
let ret = unsafe { libc::getrusage(libc::RUSAGE_THREAD, usage.as_mut_ptr()) };
|
||||
if ret != 0 {
|
||||
return None;
|
||||
}
|
||||
let usage = unsafe { usage.assume_init() };
|
||||
Some((
|
||||
timeval_to_secs(usage.ru_utime),
|
||||
timeval_to_secs(usage.ru_stime),
|
||||
))
|
||||
}
|
||||
|
||||
#[cfg(not(target_os = "linux"))]
|
||||
fn get_thread_rusage() -> Option<(f64, f64)> {
|
||||
None
|
||||
}
|
||||
|
||||
/// The `(user, system)` CPU seconds elapsed between `start` and `current`.
|
||||
///
|
||||
/// Guards against the clock going backwards
|
||||
/// (clamping to zero and logging, as the accounting must never go negative).
|
||||
fn cputime_delta(current: (f64, f64), start: (f64, f64)) -> (f64, f64) {
|
||||
let mut utime_delta = current.0 - start.0;
|
||||
let mut stime_delta = current.1 - start.1;
|
||||
|
||||
// sanity check
|
||||
if utime_delta < 0.0 {
|
||||
error!("utime went backwards! {} < {}", current.0, start.0);
|
||||
utime_delta = 0.0;
|
||||
}
|
||||
if stime_delta < 0.0 {
|
||||
error!("stime went backwards! {} < {}", current.1, start.1);
|
||||
stime_delta = 0.0;
|
||||
}
|
||||
|
||||
(utime_delta, stime_delta)
|
||||
}
|
||||
|
||||
/// 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.
|
||||
///
|
||||
/// The attribute surface, methods, error-message wording and abuse-detection
|
||||
/// behaviour are a compatibility contract with Python callers and subclasses
|
||||
/// (notably `BackgroundProcessLoggingContext`) — change both sides together.
|
||||
///
|
||||
/// Construction is split between `__new__` (which allocates a blank
|
||||
/// instance) and `__init__` (which does the real initialisation), matching 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.
|
||||
#[pyclass(subclass)]
|
||||
pub struct LoggingContext {
|
||||
/// Name for the context, used in logging. Stored as a Python string:
|
||||
/// `LoggingContextFilter` calls `str(context)` on every log record, and a
|
||||
/// `Py<PyString>` getter is INCREF-only where a `String` getter would
|
||||
/// allocate a fresh `str` per record. Rust-side reads only happen on cold
|
||||
/// error/debug paths (see [`Self::name_string`]).
|
||||
#[pyo3(get, set)]
|
||||
name: Py<PyString>,
|
||||
/// The homeserver name this context is associated with. Stored as a Python
|
||||
/// string for the same reason as `name` (read per log record).
|
||||
#[pyo3(get, set)]
|
||||
server_name: Py<PyString>,
|
||||
/// The `threading.get_ident()` value of the thread this context was created
|
||||
/// on (see [`get_thread_id`] for why it is not a real OS tid); activity on
|
||||
/// any other thread is an error. Settable only so tests can simulate
|
||||
/// activity on the wrong thread.
|
||||
#[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 CPU usage `(ru_utime, ru_stime)` in seconds captured when this
|
||||
/// context became active, or `None` if it is not currently active. Private
|
||||
/// (native `(f64, f64)` rather than a Python `struct_rusage`) so the switch
|
||||
/// path does no per-switch Python allocation; nothing outside this module
|
||||
/// reads it.
|
||||
usage_start: Option<(f64, f64)>,
|
||||
/// A short human-readable tag (e.g. the sync type). Initialised to `""` and
|
||||
/// treated as a `str` by everything in-tree, but `Option` so that assigning
|
||||
/// `None` (which the sentinel's `tag` reports, and which out-of-tree callers
|
||||
/// may assign) is accepted rather than raising `TypeError`.
|
||||
#[pyo3(get, set)]
|
||||
tag: Option<String>,
|
||||
/// The resources used by this context so far; mutated in place. Not
|
||||
/// exposed as an attribute: Python reads it via `get_resource_usage()`,
|
||||
/// which returns a copy.
|
||||
resource_usage: Py<ContextResourceUsage>,
|
||||
/// The context that was current when this one was created; restored on exit.
|
||||
/// `None` means the sentinel was current (or — only before `__init__` has
|
||||
/// run — that nothing has been recorded yet; both restore to the sentinel),
|
||||
/// so the getter reports `None` to Python where the old pure-Python
|
||||
/// attribute held `SENTINEL_CONTEXT`.
|
||||
#[pyo3(get, set)]
|
||||
previous_context: Option<Py<LoggingContext>>,
|
||||
/// 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: intern!(py, "").clone().unbind(),
|
||||
server_name: intern!(py, "").clone().unbind(),
|
||||
main_thread: 0,
|
||||
finished: false,
|
||||
usage_start: None,
|
||||
tag: Some(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: Bound<'_, PyString>,
|
||||
server_name: Bound<'_, PyString>,
|
||||
parent_context: Option<Py<PyAny>>,
|
||||
request: Option<Py<PyAny>>,
|
||||
) -> PyResult<()> {
|
||||
self.previous_context = current_context(py);
|
||||
|
||||
// The resource-usage tracker was already allocated (zeroed) by `__new__`,
|
||||
// which `type.__call__` runs immediately before this.
|
||||
|
||||
// The thread resource usage when the logcontext became active. None if
|
||||
// the context is not currently active.
|
||||
self.usage_start = None;
|
||||
|
||||
self.name = name.unbind();
|
||||
self.server_name = server_name.unbind();
|
||||
self.main_thread = get_thread_id();
|
||||
self.request = None;
|
||||
self.tag = Some(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 = parent.getattr("request")?.extract()?;
|
||||
// we also track the current scope
|
||||
self.scope = parent.getattr("scope")?.extract()?;
|
||||
}
|
||||
|
||||
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(())
|
||||
}
|
||||
|
||||
/// Returns the stored name object itself (INCREF-only): this runs per log
|
||||
/// record via `LoggingContextFilter`.
|
||||
fn __str__(&self, py: Python<'_>) -> Py<PyString> {
|
||||
self.name.clone_ref(py)
|
||||
}
|
||||
|
||||
/// Enter this logging context, making it the current context.
|
||||
fn __enter__<'py>(slf: Bound<'py, Self>) -> PyResult<Bound<'py, Self>> {
|
||||
let py = slf.py();
|
||||
// An owned handle rather than a held borrow: `set_current_context`
|
||||
// re-enters `slf` (borrow_mut in `start_inner`).
|
||||
let previous = slf
|
||||
.borrow()
|
||||
.previous_context
|
||||
.as_ref()
|
||||
.map(|p| p.clone_ref(py));
|
||||
|
||||
if log_enabled!(target: DEBUG_LOGGER_NAME, Level::Debug) {
|
||||
// The name is only materialised when the opt-in debug logger is
|
||||
// actually enabled: this runs on every context entry.
|
||||
debug!(
|
||||
target: DEBUG_LOGGER_NAME,
|
||||
"LoggingContext({}).__enter__",
|
||||
slf.borrow().name_string(py)
|
||||
);
|
||||
}
|
||||
|
||||
let old_context = set_current_context(py, Some(slf.clone().unbind()))?;
|
||||
|
||||
if !slots_identical(&previous, &old_context) {
|
||||
let previous_repr = slot_repr(py, &previous)?;
|
||||
let old_repr = slot_repr(py, &old_context)?;
|
||||
logcontext_error(
|
||||
py,
|
||||
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 previous = slf
|
||||
.borrow()
|
||||
.previous_context
|
||||
.as_ref()
|
||||
.map(|p| p.clone_ref(py));
|
||||
|
||||
if log_enabled!(target: DEBUG_LOGGER_NAME, Level::Debug) {
|
||||
// Match the Python `%s`: the str() of the previous context
|
||||
// (`"sentinel"` for an empty slot). Computed (along with the name)
|
||||
// only when the opt-in debug logger is actually enabled: this runs
|
||||
// on every context exit.
|
||||
let previous_str = match &previous {
|
||||
Some(p) => p.bind(py).str()?.extract::<String>()?,
|
||||
None => "sentinel".to_owned(),
|
||||
};
|
||||
debug!(
|
||||
target: DEBUG_LOGGER_NAME,
|
||||
"LoggingContext({}).__exit__ --> {previous_str}",
|
||||
slf.borrow().name_string(py)
|
||||
);
|
||||
}
|
||||
|
||||
let current = set_current_context(py, previous)?;
|
||||
|
||||
let restored_self = current.as_ref().is_some_and(|c| c.bind(py).is(&slf));
|
||||
if !restored_self {
|
||||
// Cold path: the name is only materialised for the error message.
|
||||
let name = slf.borrow().name_string(py);
|
||||
match ¤t {
|
||||
None => logcontext_error(py, format!("Expected logging context {name} was lost"))?,
|
||||
Some(current) => {
|
||||
let current_str: String = current.bind(py).str()?.extract()?;
|
||||
logcontext_error(
|
||||
py,
|
||||
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`. `rusage` is
|
||||
/// the thread CPU usage `(ru_utime, ru_stime)` at the point of switching to
|
||||
/// this context (`None` if the platform doesn't track it).
|
||||
fn start(slf: Bound<'_, Self>, rusage: Option<(f64, f64)>) -> PyResult<()> {
|
||||
Self::start_inner(&slf, rusage)
|
||||
}
|
||||
|
||||
/// 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<(f64, f64)>) -> PyResult<()> {
|
||||
Self::stop_inner(&slf, rusage)
|
||||
}
|
||||
|
||||
/// 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 (usage_start, main_thread) = {
|
||||
let this = slf.borrow();
|
||||
(this.usage_start, this.main_thread)
|
||||
};
|
||||
|
||||
// If we are on the correct thread and we're currently running then we can
|
||||
// include resource usage so far.
|
||||
if let Some(start) = usage_start {
|
||||
if get_thread_id() == main_thread {
|
||||
if let Some(current) = get_thread_rusage() {
|
||||
let (utime_delta, stime_delta) = cputime_delta(current, start);
|
||||
res.ru_utime += utime_delta;
|
||||
res.ru_stime += stime_delta;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(res)
|
||||
}
|
||||
|
||||
/// 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),
|
||||
|p| p.borrow().add_cputime(py, 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,),
|
||||
|p| p.borrow().add_database_transaction(py, 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,),
|
||||
|p| p.borrow().add_database_scheduled(py, 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,),
|
||||
|p| p.borrow().record_event_fetch(py, 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;
|
||||
}
|
||||
}
|
||||
|
||||
impl LoggingContext {
|
||||
/// The context name as an owned Rust string.
|
||||
///
|
||||
/// This copies the string data, so it is for cold error/debug paths only —
|
||||
/// the switch fast path must not allocate.
|
||||
fn name_string(&self, py: Python<'_>) -> String {
|
||||
self.name.bind(py).to_string_lossy().into_owned()
|
||||
}
|
||||
|
||||
/// Native body of the `start` pymethod. Shared with the switch fast path in
|
||||
/// [`set_current_context`], which calls this directly for a base
|
||||
/// `LoggingContext` rather than dispatching through Python. Runs the same
|
||||
/// thread-affinity and abuse checks on both paths.
|
||||
///
|
||||
/// This (like [`Self::stop_inner`]) runs on every context switch: the error
|
||||
/// branches materialise the name themselves so the fast path stays
|
||||
/// allocation-free.
|
||||
fn start_inner(slf: &Bound<'_, Self>, rusage: Option<(f64, f64)>) -> PyResult<()> {
|
||||
let py = slf.py();
|
||||
let main_thread = slf.borrow().main_thread;
|
||||
|
||||
if get_thread_id() != main_thread {
|
||||
let name = slf.borrow().name_string(py);
|
||||
logcontext_error(py, format!("Started logcontext {name} on different thread"))?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if slf.borrow().finished {
|
||||
let name = slf.borrow().name_string(py);
|
||||
logcontext_error(py, format!("Re-starting finished log context {name}"))?;
|
||||
}
|
||||
|
||||
// If we haven't already started, record the thread resource usage so far.
|
||||
if slf.borrow().usage_start.is_some() {
|
||||
let name = slf.borrow().name_string(py);
|
||||
logcontext_error(py, format!("Re-starting already-active log context {name}"))?;
|
||||
} else {
|
||||
slf.borrow_mut().usage_start = rusage;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Native body of the `stop` pymethod. Shared with the switch fast path.
|
||||
fn stop_inner(slf: &Bound<'_, Self>, rusage: Option<(f64, f64)>) -> PyResult<()> {
|
||||
let py = slf.py();
|
||||
let main_thread = slf.borrow().main_thread;
|
||||
|
||||
// `finally`-style: `usage_start` must be cleared however we exit.
|
||||
let result = (|| -> PyResult<()> {
|
||||
if get_thread_id() != main_thread {
|
||||
let name = slf.borrow().name_string(py);
|
||||
logcontext_error(py, format!("Stopped logcontext {name} on different thread"))?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// `if not rusage: return` — no rusage means this platform doesn't
|
||||
// track per-thread CPU, so there is nothing to account.
|
||||
let Some(current) = rusage else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
// Record the cpu used since we started.
|
||||
let Some(start) = slf.borrow().usage_start else {
|
||||
let name = slf.borrow().name_string(py);
|
||||
logcontext_error(
|
||||
py,
|
||||
format!("Called stop on logcontext {name} without recording a start rusage"),
|
||||
)?;
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
let (utime_delta, stime_delta) = cputime_delta(current, start);
|
||||
slf.borrow().add_cputime(py, utime_delta, stime_delta)?;
|
||||
Ok(())
|
||||
})();
|
||||
|
||||
slf.borrow_mut().usage_start = None;
|
||||
result
|
||||
}
|
||||
}
|
||||
|
||||
/// Which way a [`switch_context`] dispatch goes.
|
||||
#[derive(Clone, Copy)]
|
||||
enum SwitchDirection {
|
||||
Start,
|
||||
Stop,
|
||||
}
|
||||
|
||||
/// Dispatch a `stop`/`start` to a slot value during a switch, respecting
|
||||
/// subclass overrides.
|
||||
///
|
||||
/// The sentinel (an empty slot) is a no-op. A base `LoggingContext` takes the
|
||||
/// native accounting directly (the hot path — no Python dispatch, no
|
||||
/// `struct_rusage` allocation). A Python subclass (e.g.
|
||||
/// `BackgroundProcessLoggingContext`) goes through Python — materialising the
|
||||
/// rusage as a `(utime, stime)` tuple only here, off the hot path — so its
|
||||
/// overrides run. Both arms match on the same [`SwitchDirection`], so the
|
||||
/// native and Python paths provably dispatch the same operation.
|
||||
fn switch_context(
|
||||
py: Python<'_>,
|
||||
slot: Option<&Py<LoggingContext>>,
|
||||
direction: SwitchDirection,
|
||||
rusage: Option<(f64, f64)>,
|
||||
) -> PyResult<()> {
|
||||
let Some(ctx) = slot else {
|
||||
return Ok(());
|
||||
};
|
||||
let ctx = ctx.bind(py);
|
||||
if ctx.as_any().is_exact_instance_of::<LoggingContext>() {
|
||||
match direction {
|
||||
SwitchDirection::Start => LoggingContext::start_inner(ctx, rusage),
|
||||
SwitchDirection::Stop => LoggingContext::stop_inner(ctx, rusage),
|
||||
}
|
||||
} else {
|
||||
let method = match direction {
|
||||
SwitchDirection::Start => intern!(py, "start"),
|
||||
SwitchDirection::Stop => intern!(py, "stop"),
|
||||
};
|
||||
ctx.call_method1(method, (rusage,))?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether two slot values are the same context (or both the sentinel).
|
||||
/// `LoggingContext` defines no `__eq__`, so identity is the comparison Python
|
||||
/// callers got too.
|
||||
fn slots_identical(a: &Option<Py<LoggingContext>>, b: &Option<Py<LoggingContext>>) -> bool {
|
||||
match (a, b) {
|
||||
(None, None) => true,
|
||||
(Some(a), Some(b)) => a.is(b),
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
/// `repr()` of a slot value, for error messages (cold paths only). An empty
|
||||
/// slot renders as `None`, matching what the `previous_context` getter exposes
|
||||
/// to Python.
|
||||
fn slot_repr(py: Python<'_>, slot: &Option<Py<LoggingContext>>) -> PyResult<String> {
|
||||
match slot {
|
||||
Some(ctx) => Ok(ctx.bind(py).repr()?.extract()?),
|
||||
None => Ok("None".to_owned()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Set the current logging context, returning the context that was previously
|
||||
/// current. `None` means the sentinel, in both directions.
|
||||
///
|
||||
/// `context` must be a [`LoggingContext`] (or subclass) or `None` — anything
|
||||
/// else fails extraction with a `TypeError`, so the storage slots stay typed.
|
||||
/// This is not the Python-facing API: `synapse.logging.context.set_current_context`
|
||||
/// wraps this with the `SENTINEL_CONTEXT` <-> `None` mapping (and it, not this,
|
||||
/// rejects `None` from callers — here `None` legitimately means the sentinel).
|
||||
///
|
||||
/// Reads the thread rusage once (`getrusage` via libc), `stop`s the old context
|
||||
/// and `start`s the new one; for the common base-`LoggingContext` case the
|
||||
/// bookkeeping runs inline, with no per-switch Python dispatch or allocation.
|
||||
#[pyfunction]
|
||||
#[pyo3(signature = (context))]
|
||||
pub fn set_current_context(
|
||||
py: Python<'_>,
|
||||
context: Option<Py<LoggingContext>>,
|
||||
) -> PyResult<Option<Py<LoggingContext>>> {
|
||||
let current = current_context(py);
|
||||
|
||||
if !slots_identical(¤t, &context) {
|
||||
let rusage = get_thread_rusage();
|
||||
switch_context(py, current.as_ref(), SwitchDirection::Stop, rusage)?;
|
||||
// Raw slot write; we already hold `current`, so ignore the previous value
|
||||
// it returns. The clone_ref keeps a reference for the `start` below.
|
||||
let new_ref = context.as_ref().map(|ctx| ctx.clone_ref(py));
|
||||
swap_current_context(context);
|
||||
switch_context(py, new_ref.as_ref(), SwitchDirection::Start, rusage)?;
|
||||
}
|
||||
|
||||
Ok(current)
|
||||
}
|
||||
|
||||
/// Get the current logging context, or `None` for the sentinel.
|
||||
///
|
||||
/// Resolves this OS thread's slot. This is not the Python-facing API:
|
||||
/// `synapse.logging.context.current_context` wraps this and returns
|
||||
/// `SENTINEL_CONTEXT` instead of `None`.
|
||||
#[pyfunction]
|
||||
pub fn current_context(py: Python<'_>) -> Option<Py<LoggingContext>> {
|
||||
THREAD_LOCAL_CONTEXT.with(|slot| slot.borrow().as_ref().map(|ctx| ctx.clone_ref(py)))
|
||||
}
|
||||
|
||||
/// Set this OS thread's current logging context slot, returning the previous
|
||||
/// slot value (`None` is the sentinel, in both directions).
|
||||
///
|
||||
/// This is the raw slot write only — it does **not** do any resource-usage
|
||||
/// accounting or thread-affinity checks; [`set_current_context`] wraps this with
|
||||
/// the `getrusage` start/stop bookkeeping.
|
||||
///
|
||||
/// Crate-internal: a raw slot write that bypasses the rusage accounting and
|
||||
/// thread-affinity checks has no Python caller, so it is not exported (Python
|
||||
/// uses [`set_current_context`]).
|
||||
fn swap_current_context(context: Option<Py<LoggingContext>>) -> Option<Py<LoggingContext>> {
|
||||
THREAD_LOCAL_CONTEXT.with(|slot| std::mem::replace(&mut *slot.borrow_mut(), context))
|
||||
}
|
||||
|
||||
/// Called when registering modules with python.
|
||||
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!(set_current_context, &child_module)?)?;
|
||||
child_module.add("DEBUG_LOGGER_NAME", DEBUG_LOGGER_NAME)?;
|
||||
|
||||
m.add_submodule(&child_module)?;
|
||||
|
||||
@@ -152,3 +906,75 @@ pub fn register_module(py: Python<'_>, m: &Bound<'_, PyModule>) -> PyResult<()>
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use pyo3::types::PyString;
|
||||
|
||||
use super::*;
|
||||
|
||||
/// A minimal `LoggingContext` for slot tests, built directly (bypassing
|
||||
/// `__init__`, which would capture the current context and thread id).
|
||||
fn test_context(py: Python<'_>, name: &str) -> Py<LoggingContext> {
|
||||
Py::new(
|
||||
py,
|
||||
LoggingContext {
|
||||
name: PyString::new(py, name).unbind(),
|
||||
server_name: PyString::new(py, "test_server").unbind(),
|
||||
main_thread: 0,
|
||||
finished: false,
|
||||
usage_start: None,
|
||||
tag: Some(String::new()),
|
||||
resource_usage: Py::new(py, ContextResourceUsage::default())
|
||||
.expect("failed to allocate ContextResourceUsage"),
|
||||
previous_context: None,
|
||||
parent_context: None,
|
||||
request: None,
|
||||
scope: None,
|
||||
},
|
||||
)
|
||||
.expect("failed to allocate LoggingContext")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn thread_local_defaults_to_sentinel() {
|
||||
Python::initialize();
|
||||
Python::attach(|py| {
|
||||
// Nothing set on this (fresh) test thread → empty slot (the sentinel).
|
||||
assert!(current_context(py).is_none());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn swap_returns_previous_and_updates_thread_local() {
|
||||
Python::initialize();
|
||||
Python::attach(|py| {
|
||||
let a = test_context(py, "A");
|
||||
let b = test_context(py, "B");
|
||||
|
||||
// Swapping in A returns the previous slot (empty ⇒ sentinel) and
|
||||
// makes A current.
|
||||
let prev = swap_current_context(Some(a.clone_ref(py)));
|
||||
assert!(prev.is_none());
|
||||
assert!(current_context(py)
|
||||
.expect("expected a current context")
|
||||
.bind(py)
|
||||
.is(a.bind(py)));
|
||||
|
||||
// Swapping in B returns A.
|
||||
let prev = swap_current_context(Some(b.clone_ref(py)));
|
||||
assert!(prev
|
||||
.expect("expected previous context")
|
||||
.bind(py)
|
||||
.is(a.bind(py)));
|
||||
assert!(current_context(py)
|
||||
.expect("expected a current context")
|
||||
.bind(py)
|
||||
.is(b.bind(py)));
|
||||
|
||||
// Restore the empty slot (sentinel) so we don't leak into any other
|
||||
// test that happens to reuse this OS thread from the harness pool.
|
||||
swap_current_context(None);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
+34
-345
@@ -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 (
|
||||
@@ -42,6 +41,7 @@ from typing import (
|
||||
Literal,
|
||||
TypeVar,
|
||||
Union,
|
||||
cast,
|
||||
overload,
|
||||
)
|
||||
|
||||
@@ -53,20 +53,23 @@ from twisted.python.threadpool import ThreadPool
|
||||
|
||||
from synapse.logging.loggers import ExplicitlyConfiguredLogger
|
||||
from synapse.synapse_rust.logcontext import (
|
||||
DEBUG_LOGGER_NAME,
|
||||
# Not used in this module, but re-exported: callers import it from here.
|
||||
ContextResourceUsage, # noqa: F401
|
||||
LoggingContext,
|
||||
current_context as _rust_current_context,
|
||||
set_current_context as _rust_set_current_context,
|
||||
)
|
||||
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.
|
||||
|
||||
@@ -78,45 +81,12 @@ configuration and does not inherit the log level from the parent logger.
|
||||
# Restore the original logger class
|
||||
logging.setLoggerClass(original_logger_class)
|
||||
|
||||
try:
|
||||
import resource
|
||||
|
||||
# Python doesn't ship with a definition of RUSAGE_THREAD but it's defined
|
||||
# to be 1 on linux so we hard code it.
|
||||
RUSAGE_THREAD = 1
|
||||
|
||||
# If the system doesn't support RUSAGE_THREAD then this should throw an
|
||||
# exception.
|
||||
resource.getrusage(RUSAGE_THREAD)
|
||||
|
||||
is_thread_resource_usage_supported = True
|
||||
|
||||
def get_thread_resource_usage() -> "resource.struct_rusage | None":
|
||||
return resource.getrusage(RUSAGE_THREAD)
|
||||
|
||||
except Exception:
|
||||
# If the system doesn't support resource.getrusage(RUSAGE_THREAD) then we
|
||||
# won't track resource usage.
|
||||
is_thread_resource_usage_supported = False
|
||||
|
||||
def get_thread_resource_usage() -> "resource.struct_rusage | None":
|
||||
return None
|
||||
|
||||
|
||||
# a hook which can be set during testing to assert that we aren't abusing logcontexts.
|
||||
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:
|
||||
"""
|
||||
@@ -140,9 +110,6 @@ class ContextRequest:
|
||||
user_agent: str
|
||||
|
||||
|
||||
LoggingContextOrSentinel = Union["LoggingContext", "_Sentinel"]
|
||||
|
||||
|
||||
class _Sentinel:
|
||||
"""
|
||||
Sentinel to represent the root context
|
||||
@@ -176,10 +143,10 @@ class _Sentinel:
|
||||
def __str__(self) -> str:
|
||||
return "sentinel"
|
||||
|
||||
def start(self, rusage: "resource.struct_rusage | None") -> None:
|
||||
def start(self, rusage: "tuple[float, float] | None") -> None:
|
||||
pass
|
||||
|
||||
def stop(self, rusage: "resource.struct_rusage | None") -> None:
|
||||
def stop(self, rusage: "tuple[float, float] | None") -> None:
|
||||
pass
|
||||
|
||||
def add_database_transaction(self, duration_sec: float) -> None:
|
||||
@@ -197,285 +164,40 @@ class _Sentinel:
|
||||
|
||||
SENTINEL_CONTEXT = _Sentinel()
|
||||
|
||||
LoggingContextOrSentinel = Union[LoggingContext, _Sentinel]
|
||||
|
||||
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.
|
||||
def current_context() -> LoggingContextOrSentinel:
|
||||
"""Get the current logging context.
|
||||
|
||||
The storage lives in the Rust extension (see `rust/src/logging/context.rs`),
|
||||
which represents "no context" as `None` — mapped to `SENTINEL_CONTEXT` here
|
||||
so callers never see `None`.
|
||||
"""
|
||||
context = _rust_current_context()
|
||||
return SENTINEL_CONTEXT if context is None else context
|
||||
|
||||
__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()
|
||||
def set_current_context(context: LoggingContextOrSentinel) -> LoggingContextOrSentinel:
|
||||
"""Set the current logging context, returning the context that was
|
||||
previously current.
|
||||
|
||||
# track the resources used by this context so far
|
||||
self._resource_usage = ContextResourceUsage()
|
||||
The switch itself lives in the Rust extension: it reads the thread CPU usage
|
||||
once (`getrusage(RUSAGE_THREAD)`) and does the `stop`/`start` accounting
|
||||
natively. Rust represents the sentinel as `None`, converted in both
|
||||
directions here so callers only ever see `LoggingContextOrSentinel`.
|
||||
"""
|
||||
# everything blows up if we allow current_context to be set to None, so
|
||||
# sanity-check that now.
|
||||
if context is None:
|
||||
raise TypeError("'context' argument may not be None")
|
||||
|
||||
# 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)
|
||||
# The cast is needed because mypy cannot narrow the Union via the
|
||||
# `is SENTINEL_CONTEXT` identity check; Rust enforces the type at runtime.
|
||||
previous = _rust_set_current_context(
|
||||
None if context is SENTINEL_CONTEXT else cast(LoggingContext, context)
|
||||
)
|
||||
return SENTINEL_CONTEXT if previous is None else previous
|
||||
|
||||
|
||||
class LoggingContextFilter(logging.Filter):
|
||||
@@ -636,39 +358,6 @@ class PreserveLoggingContext:
|
||||
)
|
||||
|
||||
|
||||
_thread_local = threading.local()
|
||||
_thread_local.current_context = SENTINEL_CONTEXT
|
||||
|
||||
|
||||
def current_context() -> LoggingContextOrSentinel:
|
||||
"""Get the current logging context from thread local storage"""
|
||||
return getattr(_thread_local, "current_context", SENTINEL_CONTEXT)
|
||||
|
||||
|
||||
def set_current_context(context: LoggingContextOrSentinel) -> LoggingContextOrSentinel:
|
||||
"""Set the current logging context in thread local storage
|
||||
Args:
|
||||
context: The context to activate.
|
||||
|
||||
Returns:
|
||||
The context that was previously active
|
||||
"""
|
||||
# everything blows up if we allow current_context to be set to None, so sanity-check
|
||||
# that now.
|
||||
if context is None:
|
||||
raise TypeError("'context' argument may not be None")
|
||||
|
||||
current = current_context()
|
||||
|
||||
if current is not context:
|
||||
rusage = get_thread_resource_usage()
|
||||
current.stop(rusage)
|
||||
_thread_local.current_context = context
|
||||
context.start(rusage)
|
||||
|
||||
return current
|
||||
|
||||
|
||||
def nested_logging_context(suffix: str) -> LoggingContext:
|
||||
"""Creates a new logging context as a child of another.
|
||||
|
||||
|
||||
@@ -57,8 +57,6 @@ from synapse.metrics import SERVER_NAME_LABEL
|
||||
from synapse.metrics._types import Collector
|
||||
|
||||
if TYPE_CHECKING:
|
||||
import resource
|
||||
|
||||
# Old versions don't have `LiteralString`
|
||||
from typing_extensions import LiteralString
|
||||
|
||||
@@ -506,7 +504,7 @@ class BackgroundProcessLoggingContext(LoggingContext):
|
||||
desc=name, server_name=server_name, ctx=self
|
||||
)
|
||||
|
||||
def start(self, rusage: "resource.struct_rusage | None") -> None:
|
||||
def start(self, rusage: "tuple[float, float] | None") -> None:
|
||||
"""Log context has started running (again)."""
|
||||
|
||||
super().start(rusage)
|
||||
|
||||
@@ -10,7 +10,18 @@
|
||||
# See the GNU Affero General Public License for more details:
|
||||
# <https://www.gnu.org/licenses/agpl-3.0.html>.
|
||||
|
||||
from typing import Optional
|
||||
from types import TracebackType
|
||||
from typing import TYPE_CHECKING, Optional
|
||||
|
||||
from synapse.logging.context import ContextRequest
|
||||
|
||||
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."""
|
||||
@@ -29,3 +40,116 @@ class ContextResourceUsage:
|
||||
def __isub__(self, other: "ContextResourceUsage") -> "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
|
||||
"""
|
||||
|
||||
# The context that was current when this one was created; None means the
|
||||
# sentinel (the Rust storage cannot hold the pure-Python `SENTINEL_CONTEXT`
|
||||
# object, so this is narrower than the pure-Python attribute used to be).
|
||||
previous_context: Optional[LoggingContext]
|
||||
name: str
|
||||
server_name: str
|
||||
parent_context: "Optional[LoggingContext]"
|
||||
main_thread: int
|
||||
finished: bool
|
||||
request: Optional[ContextRequest]
|
||||
# Narrower than the runtime: the setter also accepts None (see the Rust
|
||||
# field docs for why), but all in-tree code treats this as str.
|
||||
tag: str
|
||||
scope: "Optional[_LogContextScope]"
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
name: str,
|
||||
server_name: str,
|
||||
parent_context: "Optional[LoggingContext]" = None,
|
||||
request: Optional[ContextRequest] = None,
|
||||
) -> None:
|
||||
"""
|
||||
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: The parent of the new context.
|
||||
request: Synapse Request Context object. Useful to associate all the
|
||||
logs happening to a given request.
|
||||
"""
|
||||
|
||||
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[tuple[float, float]]") -> None:
|
||||
"""Record that this logcontext is currently running.
|
||||
|
||||
Should not be called directly: use `set_current_context`.
|
||||
|
||||
Args:
|
||||
rusage: The thread CPU usage `(ru_utime, ru_stime)` at the point of
|
||||
switching to this context, or None if the platform doesn't
|
||||
track it.
|
||||
"""
|
||||
|
||||
def stop(self, rusage: "Optional[tuple[float, float]]") -> None:
|
||||
"""Record that this logcontext is no longer running.
|
||||
|
||||
Should not be called directly: use `set_current_context`.
|
||||
|
||||
Args:
|
||||
rusage: The thread CPU usage `(ru_utime, ru_stime)` at the point of
|
||||
switching away from this context, or None if the platform
|
||||
doesn't track it.
|
||||
"""
|
||||
|
||||
def get_resource_usage(self) -> ContextResourceUsage:
|
||||
"""Get the resources used by this logcontext so far.
|
||||
|
||||
Returns:
|
||||
A *copy* of the object tracking resource usage so far.
|
||||
"""
|
||||
|
||||
def add_cputime(self, utime_delta: float, stime_delta: float) -> None:
|
||||
"""Update the CPU time usage of this context (and any parents, recursively)."""
|
||||
|
||||
def add_database_transaction(self, duration_sec: float) -> None:
|
||||
"""Record the use of a database transaction and how long it took."""
|
||||
|
||||
def add_database_scheduled(self, sched_sec: float) -> None:
|
||||
"""Record a use of the database pool (the time taken to get a connection)."""
|
||||
|
||||
def record_event_fetch(self, event_count: int) -> None:
|
||||
"""Record a number of events being fetched from the db."""
|
||||
|
||||
def current_context() -> Optional[LoggingContext]:
|
||||
"""Get the current logging context, or None for the sentinel.
|
||||
|
||||
Resolves this OS thread's slot. This is not the Python-facing API:
|
||||
`synapse.logging.context.current_context` wraps this and returns
|
||||
`SENTINEL_CONTEXT` instead of `None`.
|
||||
"""
|
||||
|
||||
def set_current_context(
|
||||
context: Optional[LoggingContext],
|
||||
) -> Optional[LoggingContext]:
|
||||
"""Set the current logging context, returning the context that was previously
|
||||
current. `None` means the sentinel, in both directions.
|
||||
|
||||
Reads the thread CPU usage once via `getrusage(RUSAGE_THREAD)` and does the
|
||||
`stop`/`start` accounting natively. The annotated type is enforced: raises
|
||||
`TypeError` unless `context` is a `LoggingContext` (or subclass) or `None`.
|
||||
This is not the Python-facing API: `synapse.logging.context.set_current_context`
|
||||
wraps this with the `SENTINEL_CONTEXT` <-> `None` mapping.
|
||||
"""
|
||||
|
||||
@@ -20,7 +20,6 @@
|
||||
#
|
||||
|
||||
import logging
|
||||
import resource
|
||||
from contextlib import contextmanager
|
||||
from typing import Callable, Generator, cast
|
||||
from unittest.mock import patch
|
||||
@@ -707,10 +706,10 @@ class LoggingContextTestCase(unittest.TestCase):
|
||||
self.assertEqual(nested_context.name, "foo-bar")
|
||||
|
||||
|
||||
# A real (truthy) rusage for exercising the `start()`/`stop()` code paths that
|
||||
# only check whether an rusage is present. `RUSAGE_SELF` (not `RUSAGE_THREAD`)
|
||||
# so this works on platforms without per-thread rusage support.
|
||||
_TRUTHY_RUSAGE = resource.getrusage(resource.RUSAGE_SELF)
|
||||
# A stand-in rusage `(ru_utime, ru_stime)` for exercising the `start()`/`stop()`
|
||||
# code paths that only check whether an rusage is present, without depending on
|
||||
# `RUSAGE_THREAD` support or the real value of the thread's CPU clock.
|
||||
_TRUTHY_RUSAGE = (0.0, 0.0)
|
||||
|
||||
|
||||
@contextmanager
|
||||
@@ -730,12 +729,12 @@ class LogContextErrorMessageTestCase(unittest.TestCase):
|
||||
"""Characterization tests pinning the exact `logcontext_error` message shapes
|
||||
and the abuse-detection code paths.
|
||||
|
||||
These exist to guard against accidental drift in the switch machinery:
|
||||
downstream log scraping depends on the wording, argument order and the
|
||||
conditions that trigger each warning. Messages that interpolate a context
|
||||
via `%r` embed the object's `repr()` (id/address), so we reconstruct the
|
||||
expected string from the *same* live objects rather than hard-coding an
|
||||
address.
|
||||
These exist to guard against accidental drift in the switch machinery (which
|
||||
lives in Rust: `rust/src/logging/context.rs`): downstream log scraping
|
||||
depends on the wording, argument order and the conditions that trigger each
|
||||
warning. Messages that interpolate a context via `%r` embed the object's
|
||||
`repr()` (id/address), so we reconstruct the expected string from the *same*
|
||||
live objects rather than hard-coding an address.
|
||||
"""
|
||||
|
||||
def setUp(self) -> None:
|
||||
|
||||
Reference in New Issue
Block a user