1#![cfg(feature = "profiling")]
2
3use std::error::Error;
4use std::sync::OnceLock;
5
6#[cfg(feature = "profiling-bridge-pyroscope-rs")]
7use opentelemetry::trace::TraceContextExt;
8
9fn validate_pyroscope_endpoint(endpoint: &str) -> Result<(), Box<dyn Error>> {
13 use url::Url;
14
15 if endpoint.starts_with("unix://") {
17 return Ok(());
18 }
19
20 if endpoint.starts_with("http://") || endpoint.starts_with("https://") {
22 let url = Url::parse(endpoint)?;
23
24 if !url.username().is_empty() || url.password().is_some() {
26 return Err(format!(
27 "pyroscope endpoint must not contain userinfo; got: {endpoint} (ADR platform/0203 AC1)"
28 ).into());
29 }
30
31 let host = url.host_str().unwrap_or("");
32
33 match host {
34 "127.0.0.1" | "::1" | "[::1]" | "localhost" => Ok(()),
35 _ => Err(format!(
36 "pyroscope endpoint must target loopback (127.0.0.1, ::1, localhost, or unix socket); \
37 got: {endpoint} (ADR platform/0203 AC1)"
38 ).into()),
39 }
40 } else {
41 Err(
42 format!("pyroscope endpoint must be http://, https://, or unix://; got: {endpoint}")
43 .into(),
44 )
45 }
46}
47
48#[derive(Debug, Clone, Default)]
59pub(crate) struct ProfilingIdentity {
60 pub host_name: Option<String>,
62 pub deployment_environment: Option<String>,
64 pub service_version: Option<String>,
66}
67
68#[cfg(feature = "profiling-bridge-pyroscope-rs")]
69impl ProfilingIdentity {
70 fn tag_pairs(&self) -> Vec<(&'static str, &str)> {
76 let mut pairs = Vec::new();
77 if let Some(host) = self.host_name.as_deref().filter(|s| !s.is_empty()) {
78 pairs.push(("host_name", host));
79 }
80 if let Some(env) = self
81 .deployment_environment
82 .as_deref()
83 .filter(|s| !s.is_empty())
84 {
85 pairs.push(("deployment_environment", env));
86 }
87 if let Some(version) = self.service_version.as_deref().filter(|s| !s.is_empty()) {
88 pairs.push(("service_version", version));
89 }
90 pairs
91 }
92}
93
94pub struct ProfilingHandle {
97 #[cfg(feature = "profiling-bridge-pyroscope-rs")]
99 agent: Option<pyroscope::PyroscopeAgent<pyroscope::pyroscope::PyroscopeAgentRunning>>,
100 #[cfg(feature = "profiling-memory-jemalloc")]
108 memory_agent: Option<pyroscope::PyroscopeAgent<pyroscope::pyroscope::PyroscopeAgentRunning>>,
109}
110
111#[cfg(feature = "profiling-bridge-pyroscope-rs")]
112impl Drop for ProfilingHandle {
113 fn drop(&mut self) {
114 if let Some(agent) = self.agent.take() {
115 let _ = agent.stop();
116 }
117 #[cfg(feature = "profiling-memory-jemalloc")]
118 if let Some(agent) = self.memory_agent.take() {
119 let _ = agent.stop();
120 }
121 }
122}
123
124#[cfg(feature = "profiling-bridge-pyroscope-rs")]
125type BoxedTagFn = Box<dyn Fn(String, String) -> pyroscope::Result<()> + Send + Sync>;
126
127#[cfg(feature = "profiling-bridge-pyroscope-rs")]
130static PROFILING_TAG_FNS: OnceLock<(BoxedTagFn, BoxedTagFn)> = OnceLock::new();
131
132#[cfg(feature = "profiling-bridge-pyroscope-rs")]
137static PROFILING_STARTED: OnceLock<()> = OnceLock::new();
138
139#[cfg(feature = "profiling-bridge-pyroscope-rs")]
150pub(crate) fn start_pyroscope_bridge(
151 service_name: &str,
152 pyroscope_endpoint: &str,
153 identity: &ProfilingIdentity,
154) -> Result<Option<ProfilingHandle>, Box<dyn Error>> {
155 use pyroscope::backend::{BackendConfig, PprofConfig, pprof_backend};
156
157 validate_pyroscope_endpoint(pyroscope_endpoint)?;
159
160 if PROFILING_STARTED.set(()).is_err() {
163 return Ok(None);
164 }
165
166 let tags = identity.tag_pairs();
167
168 let agent = pyroscope::pyroscope::PyroscopeAgentBuilder::new(
169 pyroscope_endpoint,
170 service_name,
171 100,
172 "pyroscope-rs",
173 env!("CARGO_PKG_VERSION"),
174 pprof_backend(PprofConfig { sample_rate: 100 }, BackendConfig::default()),
175 )
176 .tags(tags.clone())
177 .build()?
178 .start()?;
179
180 let (add_tag, remove_tag) = agent.tag_wrapper();
181 PROFILING_TAG_FNS
182 .set((Box::new(add_tag), Box::new(remove_tag)))
183 .ok();
184
185 Ok(Some(ProfilingHandle {
186 agent: Some(agent),
187 #[cfg(feature = "profiling-memory-jemalloc")]
188 memory_agent: start_memory_agent(service_name, pyroscope_endpoint, &tags)?,
189 }))
190}
191
192#[cfg(feature = "profiling-memory-jemalloc")]
202fn start_memory_agent(
203 service_name: &str,
204 pyroscope_endpoint: &str,
205 tags: &[(&'static str, &str)],
206) -> Result<
207 Option<pyroscope::PyroscopeAgent<pyroscope::pyroscope::PyroscopeAgentRunning>>,
208 Box<dyn Error>,
209> {
210 use pyroscope::backend::jemalloc::jemalloc_backend;
211
212 let built = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
223 pyroscope::pyroscope::PyroscopeAgentBuilder::new(
224 pyroscope_endpoint,
225 service_name,
226 100,
227 "pyroscope-rs",
228 env!("CARGO_PKG_VERSION"),
229 jemalloc_backend(),
230 )
231 .tags(tags.to_vec())
232 .build()
233 }));
234
235 let agent = match built {
236 Ok(Ok(agent)) => agent,
237 Ok(Err(e)) => {
238 tracing::warn!(
239 error = %e,
240 "jemalloc heap profiling unavailable — continuing without it; \
241 check the global allocator is jemalloc and prof:true,prof_active:true is set"
242 );
243 return Ok(None);
244 }
245 Err(_) => {
246 tracing::warn!(
247 "jemalloc heap profiling unavailable — this process is not using \
248 jemalloc as its global allocator; continuing without it"
249 );
250 return Ok(None);
251 }
252 };
253
254 match agent.start() {
255 Ok(running) => {
256 tracing::info!("jemalloc heap profiling started");
257 Ok(Some(running))
258 }
259 Err(e) => {
260 tracing::warn!(error = %e, "jemalloc heap profiling failed to start — continuing without it");
261 Ok(None)
262 }
263 }
264}
265
266#[cfg(all(feature = "profiling", not(feature = "profiling-bridge-pyroscope-rs")))]
268pub(crate) fn start_pyroscope_bridge(
269 _service_name: &str,
270 _pyroscope_endpoint: &str,
271 _identity: &ProfilingIdentity,
272) -> Result<Option<ProfilingHandle>, Box<dyn Error>> {
273 Ok(None)
274}
275
276#[cfg(feature = "profiling-bridge-pyroscope-rs")]
279pub struct ProfilingTagLayer;
280
281#[cfg(feature = "profiling-bridge-pyroscope-rs")]
282impl<S> tracing_subscriber::Layer<S> for ProfilingTagLayer
283where
284 S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
285{
286 fn on_enter(&self, _id: &tracing::span::Id, _ctx: tracing_subscriber::layer::Context<'_, S>) {
287 if let Some((add_tag, _)) = PROFILING_TAG_FNS.get() {
288 let cx = opentelemetry::Context::current();
289 let span_ref = cx.span();
290 let span_context = span_ref.span_context();
291 if span_context.is_valid() {
292 let trace_id = span_context.trace_id();
293 let span_id = span_context.span_id();
294 let _ = add_tag("trace_id".to_string(), format!("{trace_id:x}"));
295 let _ = add_tag("span_id".to_string(), format!("{span_id:x}"));
296 }
297 }
298 }
299
300 fn on_exit(&self, _id: &tracing::span::Id, _ctx: tracing_subscriber::layer::Context<'_, S>) {
301 if let Some((_, remove_tag)) = PROFILING_TAG_FNS.get() {
302 let cx = opentelemetry::Context::current();
303 let span_ref = cx.span();
304 let span_context = span_ref.span_context();
305 if span_context.is_valid() {
306 let trace_id = span_context.trace_id();
307 let span_id = span_context.span_id();
308 let _ = remove_tag("trace_id".to_string(), format!("{trace_id:x}"));
309 let _ = remove_tag("span_id".to_string(), format!("{span_id:x}"));
310 }
311 }
312 }
313}
314
315#[cfg(all(test, feature = "profiling-bridge-pyroscope-rs"))]
316mod tests {
317 use super::*;
318
319 #[test]
320 fn start_bridge_with_nonexistent_server() {
321 let result = start_pyroscope_bridge(
322 "test-svc",
323 "http://localhost:4040",
324 &ProfilingIdentity::default(),
325 );
326 assert!(
327 result.is_ok(),
328 "pyroscope agent start() is lazy and does not eagerly connect"
329 );
330 if let Ok(Some(_handle)) = result {
331 }
333 }
334
335 #[test]
336 fn start_bridge_multiple_times_ignores_second() {
337 let result1 = start_pyroscope_bridge(
338 "test-svc-1",
339 "http://localhost:4040",
340 &ProfilingIdentity::default(),
341 );
342 assert!(result1.is_ok());
343 let result2 = start_pyroscope_bridge(
344 "test-svc-2",
345 "http://localhost:4041",
346 &ProfilingIdentity::default(),
347 );
348 assert!(result2.is_ok());
349 assert!(result2.unwrap().is_none());
352 }
353
354 #[test]
355 fn validate_endpoint_accepts_loopback_ipv4() {
356 assert!(validate_pyroscope_endpoint("http://127.0.0.1:4040").is_ok());
357 }
358
359 #[test]
360 fn validate_endpoint_accepts_loopback_ipv6() {
361 assert!(validate_pyroscope_endpoint("http://[::1]:4040").is_ok());
363 }
364
365 #[test]
366 fn validate_endpoint_accepts_localhost() {
367 assert!(validate_pyroscope_endpoint("http://localhost:4040").is_ok());
368 }
369
370 #[test]
371 fn validate_endpoint_accepts_https_loopback() {
372 assert!(validate_pyroscope_endpoint("https://127.0.0.1:4040").is_ok());
373 }
374
375 #[test]
376 fn validate_endpoint_rejects_routable_ipv4() {
377 assert!(validate_pyroscope_endpoint("http://10.0.0.1:4040").is_err());
378 }
379
380 #[test]
381 fn validate_endpoint_rejects_userinfo_bypass() {
382 assert!(validate_pyroscope_endpoint("http://127.0.0.1:4040@evil.com/").is_err());
384 }
385
386 #[test]
387 fn validate_endpoint_rejects_userinfo_with_password() {
388 assert!(validate_pyroscope_endpoint("http://user:pass@localhost:4040").is_err());
389 }
390
391 #[test]
392 fn validate_endpoint_rejects_unix_socket_check() {
393 assert!(validate_pyroscope_endpoint("unix:///var/run/profiling.sock").is_ok());
394 }
395}
396
397#[cfg(all(
398 test,
399 feature = "profiling",
400 not(feature = "profiling-bridge-pyroscope-rs")
401))]
402mod tests_no_bridge {
403 use super::*;
404
405 #[test]
406 fn start_bridge_returns_none() {
407 let result = start_pyroscope_bridge(
408 "test-svc",
409 "http://localhost:4040",
410 &ProfilingIdentity::default(),
411 );
412 assert!(result.is_ok());
413 if let Ok(handle) = result {
414 assert!(handle.is_none());
415 }
416 }
417}