mirror of
https://github.com/XRPLF/rippled.git
synced 2026-07-23 23:20:33 +00:00
Merge branch 'pratik/otel-phase1b-telemetry-infra' into pratik/otel-phase1c-rpc-integration
Brings coroutine-aware OTel context storage forward: CoroAwareContextStorage, its install in Telemetry, ScopedSpanGuard same-store assertion, and the non-owning ScopedActivation helper. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
68
src/libxrpl/telemetry/CoroAwareContextStorage.cpp
Normal file
68
src/libxrpl/telemetry/CoroAwareContextStorage.cpp
Normal file
@@ -0,0 +1,68 @@
|
||||
/**
|
||||
* Implementation of CoroAwareContextStorage.
|
||||
*
|
||||
* Mirrors opentelemetry's ThreadLocalContextStorage stack semantics, but the
|
||||
* stack lives in an xrpl::LocalValue so it follows a JobQueue::Coro across
|
||||
* yield/resume. All OpenTelemetry types stay confined to this translation unit.
|
||||
*
|
||||
* @see CoroAwareContextStorage (CoroAwareContextStorage.h)
|
||||
*/
|
||||
|
||||
#ifdef XRPL_ENABLE_TELEMETRY
|
||||
|
||||
#include <xrpl/telemetry/CoroAwareContextStorage.h>
|
||||
|
||||
#include <opentelemetry/context/context.h>
|
||||
#include <opentelemetry/context/runtime_context.h>
|
||||
#include <opentelemetry/nostd/unique_ptr.h>
|
||||
|
||||
namespace xrpl::telemetry {
|
||||
|
||||
opentelemetry::context::Context
|
||||
CoroAwareContextStorage::GetCurrent() noexcept
|
||||
{
|
||||
auto& stack = *stack_;
|
||||
return stack.empty() ? opentelemetry::context::Context{} : stack.back();
|
||||
}
|
||||
|
||||
opentelemetry::nostd::unique_ptr<opentelemetry::context::Token>
|
||||
CoroAwareContextStorage::Attach(opentelemetry::context::Context const& context) noexcept
|
||||
{
|
||||
stack_->push_back(context);
|
||||
return CreateToken(context);
|
||||
}
|
||||
|
||||
bool
|
||||
CoroAwareContextStorage::Detach(opentelemetry::context::Token& token) noexcept
|
||||
{
|
||||
auto& stack = *stack_;
|
||||
// Fast path: the token's frame is on top (the common LIFO case).
|
||||
if (!stack.empty() && token == stack.back())
|
||||
{
|
||||
stack.pop_back();
|
||||
return true;
|
||||
}
|
||||
// Fallback: token not on top — verify it is present, then pop down to and
|
||||
// including it (mirrors ThreadLocalContextStorage::Detach; also detaches
|
||||
// any child frames left above it).
|
||||
bool found = false;
|
||||
for (auto const& frame : stack)
|
||||
{
|
||||
if (token == frame)
|
||||
{
|
||||
found = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (!found)
|
||||
return false;
|
||||
while (!stack.empty() && !(token == stack.back()))
|
||||
stack.pop_back();
|
||||
if (!stack.empty())
|
||||
stack.pop_back();
|
||||
return true;
|
||||
}
|
||||
|
||||
} // namespace xrpl::telemetry
|
||||
|
||||
#endif // XRPL_ENABLE_TELEMETRY
|
||||
@@ -10,10 +10,12 @@
|
||||
* touches the thread-local context stack, so a SpanGuard is
|
||||
* thread-free and may be moved to and destroyed on any thread.
|
||||
* - ScopedSpanGuard::ScopedImpl holds a SpanGuard plus an optional
|
||||
* Scope. The Scope pushes the span onto the constructing thread's
|
||||
* context stack; member order (guard first, scope second) ensures
|
||||
* Scope. The Scope pushes the span onto the constructing context
|
||||
* store's stack; member order (guard first, scope second) ensures
|
||||
* the Scope pops BEFORE the span ends on destruction. It is
|
||||
* thread-bound: construct and destroy on the same thread.
|
||||
* store-bound: construct and destroy while the same LocalValue
|
||||
* context store is active (the same thread, or the same JobQueue
|
||||
* coroutine even if it resumes on another worker).
|
||||
*
|
||||
* Static factory methods access the global Telemetry instance via
|
||||
* Telemetry::getInstance(), check whether the requested TraceCategory
|
||||
@@ -28,6 +30,7 @@
|
||||
|
||||
#include <xrpl/telemetry/SpanGuard.h>
|
||||
|
||||
#include <xrpl/basics/LocalValue.h>
|
||||
#include <xrpl/beast/utility/instrumentation.h>
|
||||
#include <xrpl/telemetry/DiscardFlag.h>
|
||||
#include <xrpl/telemetry/Telemetry.h>
|
||||
@@ -48,7 +51,6 @@
|
||||
#include <optional>
|
||||
#include <string>
|
||||
#include <string_view>
|
||||
#include <thread>
|
||||
#include <typeinfo>
|
||||
#include <utility>
|
||||
|
||||
@@ -432,21 +434,35 @@ struct ScopedSpanGuard::ScopedImpl
|
||||
SpanGuard guard;
|
||||
|
||||
/**
|
||||
* Active OTel Scope binding the span to this thread's context stack.
|
||||
* Present while the guard is active; reset() (via operator SpanGuard)
|
||||
* pops it eagerly on the origin thread. Declared AFTER `guard` so
|
||||
* destruction order pops the scope before the span ends.
|
||||
* Active OTel Scope binding the span to the active context store's
|
||||
* stack. Present while the guard is active; reset() (via operator
|
||||
* SpanGuard) pops it eagerly under the constructing store. Declared
|
||||
* AFTER `guard` so destruction order pops the scope before the span
|
||||
* ends.
|
||||
*/
|
||||
std::optional<otel_trace::Scope> scope;
|
||||
|
||||
/**
|
||||
* Thread that constructed this ScopedImpl. The Scope is bound to
|
||||
* this thread's context stack, so popping it (via ~Scope or reset())
|
||||
* on a different thread corrupts that thread's stack. Checked in the
|
||||
* destructor and operator SpanGuard() to turn a silent cross-thread
|
||||
* corruption into an assertion failure in debug/test/fuzzing builds.
|
||||
* Identity of the LocalValue store the Scope was pushed onto (the
|
||||
* coroutine's store on a coro, else the thread's own store). Popping
|
||||
* the Scope (via ~Scope or reset()) must happen while the SAME store
|
||||
* is active, else a different stack is corrupted. Coroutine-aware
|
||||
* storage lets a coro legitimately resume on another worker while
|
||||
* keeping the same store, so store-identity — not thread-id — is the
|
||||
* correct invariant. Checked in the destructor and operator
|
||||
* SpanGuard() to turn a silent cross-store pop into an assertion
|
||||
* failure in debug/test/fuzzing builds.
|
||||
*
|
||||
* Assigned in the constructor body AFTER the Scope push, not as a
|
||||
* member initializer: the push is the first LocalValue touch on a
|
||||
* fresh thread for a root or explicit-parent span (the OTel SDK skips
|
||||
* the ambient-context read in those cases), so it materializes the
|
||||
* store. Capturing before the push would record nullptr while the
|
||||
* destructor sees the now-materialized store. For a null guard no
|
||||
* Scope is pushed and the assertions short-circuit, so this holds
|
||||
* whatever store is active then.
|
||||
*/
|
||||
std::thread::id owner{std::this_thread::get_id()};
|
||||
void const* owner = nullptr;
|
||||
|
||||
/**
|
||||
* Wrap a SpanGuard and activate its span on this thread. If the
|
||||
@@ -458,6 +474,9 @@ struct ScopedSpanGuard::ScopedImpl
|
||||
{
|
||||
if (guard)
|
||||
scope.emplace(guard.impl_->span);
|
||||
// Capture after the push so `owner` names the store the Scope
|
||||
// actually lives on, even when the push materialized it.
|
||||
owner = detail::getLocalValues().get();
|
||||
}
|
||||
};
|
||||
|
||||
@@ -470,13 +489,14 @@ ScopedSpanGuard::ScopedSpanGuard(SpanGuard&& guard)
|
||||
|
||||
ScopedSpanGuard::~ScopedSpanGuard()
|
||||
{
|
||||
// A live Scope must be popped on the thread that pushed it; destroying
|
||||
// the guard on any other thread corrupts that thread's context stack.
|
||||
// A null guard holds no Scope, so the check is skipped.
|
||||
// A live Scope must be popped while its constructing context store is
|
||||
// active; destroying the guard under a different store corrupts that
|
||||
// store's context stack. A null guard holds no Scope, so the check is
|
||||
// skipped.
|
||||
XRPL_ASSERT(
|
||||
!impl_ || !impl_->scope.has_value() || impl_->owner == std::this_thread::get_id(),
|
||||
"xrpl::telemetry::ScopedSpanGuard::~ScopedSpanGuard : destroyed on "
|
||||
"constructing thread");
|
||||
!impl_ || !impl_->scope.has_value() || impl_->owner == detail::getLocalValues().get(),
|
||||
"xrpl::telemetry::ScopedSpanGuard::~ScopedSpanGuard : destroyed on the "
|
||||
"constructing context store");
|
||||
}
|
||||
|
||||
// ===== ScopedSpanGuard factory methods =====================================
|
||||
@@ -521,13 +541,14 @@ ScopedSpanGuard::linkedSpan(std::string_view name, SpanContext const& linkCtx)
|
||||
ScopedSpanGuard::
|
||||
operator SpanGuard() &&
|
||||
{
|
||||
// The scope must be popped on the thread that pushed it; handing off
|
||||
// elsewhere would pop the wrong stack.
|
||||
// The scope must be popped while its constructing context store is
|
||||
// active; handing off under a different store would pop the wrong stack.
|
||||
// A null guard holds no Scope, so the check is skipped.
|
||||
XRPL_ASSERT(
|
||||
impl_->owner == std::this_thread::get_id(),
|
||||
"xrpl::telemetry::ScopedSpanGuard::operator SpanGuard : handoff on "
|
||||
"constructing thread");
|
||||
impl_->scope.reset(); // eager pop, on the origin thread
|
||||
!impl_->scope.has_value() || impl_->owner == detail::getLocalValues().get(),
|
||||
"xrpl::telemetry::ScopedSpanGuard::operator SpanGuard : handoff on the "
|
||||
"constructing context store");
|
||||
impl_->scope.reset(); // eager pop, on the origin store
|
||||
return std::move(impl_->guard);
|
||||
}
|
||||
|
||||
@@ -598,14 +619,15 @@ ScopedSpanGuard::recordException(std::exception const& e)
|
||||
void
|
||||
ScopedSpanGuard::discard()
|
||||
{
|
||||
// A live Scope must be popped on the thread that pushed it; discarding
|
||||
// on any other thread corrupts that thread's context stack. A null guard
|
||||
// holds no Scope, so the check is skipped.
|
||||
// A live Scope must be popped while its constructing context store is
|
||||
// active; discarding under a different store corrupts that store's
|
||||
// context stack. A null guard holds no Scope, so the check is skipped.
|
||||
XRPL_ASSERT(
|
||||
!impl_ || !impl_->scope.has_value() || impl_->owner == std::this_thread::get_id(),
|
||||
"xrpl::telemetry::ScopedSpanGuard::discard : discarded on constructing thread");
|
||||
// Pop the scope first (on this thread) so no span stays active on the
|
||||
// stack, then discard the span via the owned guard.
|
||||
!impl_ || !impl_->scope.has_value() || impl_->owner == detail::getLocalValues().get(),
|
||||
"xrpl::telemetry::ScopedSpanGuard::discard : discarded on the "
|
||||
"constructing context store");
|
||||
// Pop the scope first (under the constructing store) so no span stays
|
||||
// active on the stack, then discard the span via the owned guard.
|
||||
impl_->scope.reset();
|
||||
impl_->guard.discard();
|
||||
}
|
||||
@@ -616,6 +638,64 @@ operator bool() const
|
||||
return impl_ && static_cast<bool>(impl_->guard);
|
||||
}
|
||||
|
||||
// ===== ScopedActivation ====================================================
|
||||
|
||||
struct ScopedActivation::Impl
|
||||
{
|
||||
/**
|
||||
* Active OTel Scope binding an externally-owned span to the current
|
||||
* context store. Popped on destruction; never ends the span.
|
||||
*/
|
||||
otel_trace::Scope scope;
|
||||
|
||||
/**
|
||||
* Identity of the LocalValue store the Scope was pushed onto. The Scope
|
||||
* must pop while the SAME store is active (see
|
||||
* ScopedSpanGuard::ScopedImpl::owner). Initialized AFTER `scope` in the
|
||||
* member init list: members initialize in declaration order, so the
|
||||
* Scope's push (the first LocalValue touch on a fresh worker, which
|
||||
* materializes the store) runs first, then this captures the
|
||||
* now-materialized store. Capturing before the push would record
|
||||
* nullptr while the destructor sees the materialized store.
|
||||
*/
|
||||
void const* owner;
|
||||
|
||||
/**
|
||||
* Push an externally-owned span onto the current context store.
|
||||
* @param span The already-owned span to activate (not owned here).
|
||||
*/
|
||||
explicit Impl(opentelemetry::nostd::shared_ptr<otel_trace::Span> const& span)
|
||||
: scope(span), owner(detail::getLocalValues().get())
|
||||
{
|
||||
}
|
||||
};
|
||||
|
||||
ScopedActivation::ScopedActivation() = default;
|
||||
|
||||
ScopedActivation::ScopedActivation(std::unique_ptr<Impl> impl) : impl_(std::move(impl))
|
||||
{
|
||||
}
|
||||
|
||||
ScopedActivation::~ScopedActivation()
|
||||
{
|
||||
// A live Scope must be popped while its constructing context store is
|
||||
// active; destroying the activation under a different store corrupts
|
||||
// that store's context stack. A null activation holds no Scope, so the
|
||||
// check is skipped.
|
||||
XRPL_ASSERT(
|
||||
!impl_ || impl_->owner == detail::getLocalValues().get(),
|
||||
"xrpl::telemetry::ScopedActivation::~ScopedActivation : destroyed on the "
|
||||
"constructing context store");
|
||||
}
|
||||
|
||||
ScopedActivation
|
||||
SpanGuard::activate() const
|
||||
{
|
||||
if (!impl_ || !impl_->span)
|
||||
return {};
|
||||
return ScopedActivation(std::make_unique<ScopedActivation::Impl>(impl_->span));
|
||||
}
|
||||
|
||||
} // namespace xrpl::telemetry
|
||||
|
||||
#endif // XRPL_ENABLE_TELEMETRY
|
||||
|
||||
@@ -20,11 +20,13 @@
|
||||
|
||||
#include <xrpl/basics/Log.h>
|
||||
#include <xrpl/beast/utility/Journal.h>
|
||||
#include <xrpl/telemetry/CoroAwareContextStorage.h>
|
||||
#include <xrpl/telemetry/DeterministicIdGenerator.h>
|
||||
#include <xrpl/telemetry/DiscardFlag.h>
|
||||
#include <xrpl/telemetry/SpanNames.h>
|
||||
|
||||
#include <opentelemetry/context/context.h>
|
||||
#include <opentelemetry/context/runtime_context.h>
|
||||
#include <opentelemetry/exporters/otlp/otlp_http_exporter_factory.h>
|
||||
#include <opentelemetry/exporters/otlp/otlp_http_exporter_options.h>
|
||||
#include <opentelemetry/nostd/shared_ptr.h>
|
||||
@@ -264,6 +266,13 @@ class TelemetryImpl : public Telemetry
|
||||
*/
|
||||
std::shared_ptr<trace_sdk::TracerProvider> sdkProvider_;
|
||||
|
||||
/**
|
||||
* Coroutine-aware runtime-context storage, installed globally so the OTel
|
||||
* ambient context follows JobQueue coroutines. Held for the process
|
||||
* lifetime because it must outlive every span (SDK requirement).
|
||||
*/
|
||||
opentelemetry::nostd::shared_ptr<opentelemetry::context::RuntimeContextStorage> contextStorage_;
|
||||
|
||||
public:
|
||||
TelemetryImpl(Setup setup, beast::Journal journal) : setup_(std::move(setup)), journal_(journal)
|
||||
{
|
||||
@@ -334,6 +343,18 @@ public:
|
||||
std::move(sampler),
|
||||
std::make_unique<DeterministicIdGenerator>());
|
||||
|
||||
// Install coroutine-aware context storage BEFORE any span is created
|
||||
// so the OTel ambient context follows JobQueue coroutines across
|
||||
// yield/resume (fixes wrong-thread scope pop; keeps log-trace
|
||||
// correlation). Must precede SetTracerProvider and the first span.
|
||||
// Not reset in stop(): resetting the storage while spans may still
|
||||
// exist is undefined behaviour (SDK), and by stop() all spans are
|
||||
// gone, so the storage is simply left installed for process lifetime.
|
||||
contextStorage_ =
|
||||
opentelemetry::nostd::shared_ptr<opentelemetry::context::RuntimeContextStorage>(
|
||||
new CoroAwareContextStorage());
|
||||
opentelemetry::context::RuntimeContext::SetRuntimeContextStorage(contextStorage_);
|
||||
|
||||
// Set as global provider
|
||||
trace_api::Provider::SetTracerProvider(
|
||||
opentelemetry::nostd::shared_ptr<trace_api::TracerProvider>(sdkProvider_));
|
||||
|
||||
Reference in New Issue
Block a user