11use arc_swap:: ArcSwapOption ;
2- use protocol:: Event ;
2+ use protocol:: { Args , Event , TimestampType } ;
33use std:: sync:: atomic:: { AtomicBool , Ordering } ;
44use std:: sync:: Arc ;
55use std:: thread:: { self , JoinHandle } ;
66use thread_local:: ThreadLocal ;
7+ use tracing_subscriber:: EnvFilter ;
78
89use crate :: epoll_thread:: epoll_listener_thread;
910use crate :: { Producer , Result } ;
@@ -12,14 +13,35 @@ use crate::{Producer, Result};
1213pub ( crate ) struct ProducerState {
1314 producer : Arc < Producer > ,
1415 clock_id : libc:: clockid_t ,
16+ env_filter : Option < EnvFilter > ,
1517 thread_names_sent : Arc < ThreadLocal < std:: cell:: Cell < bool > > > ,
1618}
1719
1820impl ProducerState {
19- pub ( crate ) fn new ( producer : Arc < Producer > , clock_id : libc:: clockid_t ) -> Self {
21+ pub ( crate ) fn new ( producer : Arc < Producer > , args : Option < Args > ) -> Self {
22+ let ( clock_id, env_filter) = match & args {
23+ Some ( Args :: Tracing ( tracing_args) ) => {
24+ let clock_id = match tracing_args. timestamp_type {
25+ TimestampType :: Monotonic => libc:: CLOCK_MONOTONIC ,
26+ TimestampType :: Boottime => libc:: CLOCK_BOOTTIME ,
27+ TimestampType :: Realtime => libc:: CLOCK_REALTIME ,
28+ } ;
29+ // filter has already been validated in agent_state.rs
30+ let filter = EnvFilter :: try_new ( & tracing_args. log_filter ) . unwrap ( ) ;
31+ tracing:: debug!(
32+ "computed clock_id: {}, log_filter: {}" ,
33+ clock_id,
34+ tracing_args. log_filter
35+ ) ;
36+ ( clock_id, Some ( filter) )
37+ }
38+ None => ( libc:: CLOCK_MONOTONIC , None ) ,
39+ } ;
40+
2041 Self {
2142 producer,
2243 clock_id,
44+ env_filter,
2345 thread_names_sent : Arc :: new ( ThreadLocal :: new ( ) ) ,
2446 }
2547 }
@@ -28,6 +50,10 @@ impl ProducerState {
2850 self . clock_id
2951 }
3052
53+ pub ( crate ) fn env_filter ( & self ) -> Option < & EnvFilter > {
54+ self . env_filter . as_ref ( )
55+ }
56+
3157 pub ( crate ) fn submit_with_thread_name ( & self , event : & Event ) -> Result < ( ) > {
3258 let thread_sent = self
3359 . thread_names_sent
@@ -119,6 +145,15 @@ impl Agent {
119145 . unwrap_or ( libc:: CLOCK_MONOTONIC )
120146 }
121147
148+ /// Get the current environment filter if a client is connected.
149+ pub fn env_filter < F , R > ( & self , f : F ) -> Option < R >
150+ where
151+ F : FnOnce ( & EnvFilter ) -> R ,
152+ {
153+ let guard = self . producer_state . load ( ) ;
154+ guard. as_ref ( ) . and_then ( |s| s. env_filter ( ) ) . map ( f)
155+ }
156+
122157 /// Gracefully shut down the agent and wait for thread termination.
123158 pub fn wait_terminated ( & mut self ) -> Result < ( ) > {
124159 self . shutdown . store ( true , Ordering :: Relaxed ) ;
@@ -142,7 +177,7 @@ impl Drop for Agent {
142177mod tests {
143178 use super :: * ;
144179 use crate :: { AgentClient , Consumer } ;
145- use protocol:: { Counter , Event , Labels , LogLevel } ;
180+ use protocol:: { Counter , Event , Labels } ;
146181 use rstest:: * ;
147182 use std:: borrow:: Cow ;
148183 use std:: sync:: Once ;
@@ -213,7 +248,7 @@ mod tests {
213248 wait_for_socket ( & temp_socket_path) ;
214249
215250 let mut client = AgentClient :: new ( temp_socket_path) ;
216- client. start ( & consumer, LogLevel :: Debug ) . unwrap ( ) ;
251+ client. start ( & consumer, "debug" . to_string ( ) ) . unwrap ( ) ;
217252 assert ! ( client. enabled( ) ) ;
218253 }
219254
@@ -225,11 +260,11 @@ mod tests {
225260 wait_for_socket ( & temp_socket_path) ;
226261
227262 let mut client1 = AgentClient :: new ( temp_socket_path. clone ( ) ) ;
228- client1. start ( & consumer, LogLevel :: Debug ) . unwrap ( ) ;
263+ client1. start ( & consumer, "debug" . to_string ( ) ) . unwrap ( ) ;
229264 assert ! ( client1. enabled( ) ) ;
230265
231266 let mut client2 = AgentClient :: new ( temp_socket_path) ;
232- let result = client2. start ( & consumer, LogLevel :: Debug ) ;
267+ let result = client2. start ( & consumer, "debug" . to_string ( ) ) ;
233268 assert ! ( result. is_err( ) ) ;
234269 assert ! ( !client2. enabled( ) ) ;
235270 }
@@ -242,7 +277,7 @@ mod tests {
242277 wait_for_socket ( & temp_socket_path) ;
243278
244279 let mut client = AgentClient :: new ( temp_socket_path) ;
245- client. start ( & consumer, LogLevel :: Debug ) . unwrap ( ) ;
280+ client. start ( & consumer, "debug" . to_string ( ) ) . unwrap ( ) ;
246281 assert ! ( client. enabled( ) ) ;
247282
248283 client. send_continue ( ) . unwrap ( ) ;
@@ -258,7 +293,7 @@ mod tests {
258293 wait_for_socket ( & temp_socket_path) ;
259294
260295 let mut client = AgentClient :: new ( temp_socket_path) ;
261- client. start ( & consumer, LogLevel :: Debug ) . unwrap ( ) ;
296+ client. start ( & consumer, "debug" . to_string ( ) ) . unwrap ( ) ;
262297
263298 thread:: sleep ( Duration :: from_millis ( 10 ) ) ;
264299
@@ -284,7 +319,7 @@ mod tests {
284319 wait_for_socket ( & temp_socket_path) ;
285320
286321 let mut client = AgentClient :: new ( temp_socket_path) ;
287- client. start ( & consumer, LogLevel :: Debug ) . unwrap ( ) ;
322+ client. start ( & consumer, "debug" . to_string ( ) ) . unwrap ( ) ;
288323
289324 thread:: sleep ( Duration :: from_millis ( 10 ) ) ;
290325
@@ -314,7 +349,7 @@ mod tests {
314349 wait_for_socket ( & temp_socket_path) ;
315350
316351 let mut client = AgentClient :: new ( temp_socket_path. clone ( ) ) ;
317- client. start ( & consumer, LogLevel :: Debug ) . unwrap ( ) ;
352+ client. start ( & consumer, "debug" . to_string ( ) ) . unwrap ( ) ;
318353
319354 thread:: sleep ( Duration :: from_millis ( 10 ) ) ;
320355
0 commit comments