Sitelet https://github.com/category-labs/manytrace/commit/c7efcc018c6be8ac39e27018bc9506601e733330
Skip to content

Commit c7efcc0

Browse files
committed
agent, protocol, tracing-manytrace: implement dynamic log filtering with Args enum
Replace LogLevel enum with string-based log filters and add EnvFilter validation. Store computed clock_id and EnvFilter in ProducerState for efficient filtering.
1 parent 4282a0f commit c7efcc0

12 files changed

Lines changed: 165 additions & 63 deletions

File tree

‎agent/Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ nix = { workspace = true, features = ["socket", "uio", "event"] }
1313
rkyv.workspace = true
1414
thiserror.workspace = true
1515
tracing.workspace = true
16+
tracing-subscriber.workspace = true
1617
libc.workspace = true
1718
thread_local = "1.1"
1819

‎agent/src/agent.rs‎

Lines changed: 45 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,10 @@
11
use arc_swap::ArcSwapOption;
2-
use protocol::Event;
2+
use protocol::{Args, Event, TimestampType};
33
use std::sync::atomic::{AtomicBool, Ordering};
44
use std::sync::Arc;
55
use std::thread::{self, JoinHandle};
66
use thread_local::ThreadLocal;
7+
use tracing_subscriber::EnvFilter;
78

89
use crate::epoll_thread::epoll_listener_thread;
910
use crate::{Producer, Result};
@@ -12,14 +13,35 @@ use crate::{Producer, Result};
1213
pub(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

1820
impl 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 {
142177
mod 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

‎agent/src/agent_state.rs‎

Lines changed: 37 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -11,14 +11,6 @@ use tracing::{debug, warn};
1111
use crate::agent::ProducerState;
1212
use crate::{AgentError, Producer, Result};
1313

14-
fn timestamp_type_to_clock_id(timestamp_type: TimestampType) -> libc::clockid_t {
15-
match timestamp_type {
16-
TimestampType::Monotonic => libc::CLOCK_MONOTONIC,
17-
TimestampType::Boottime => libc::CLOCK_BOOTTIME,
18-
TimestampType::Realtime => libc::CLOCK_REALTIME,
19-
}
20-
}
21-
2214
pub(crate) struct AgentContext {
2315
pub(crate) another_started: bool,
2416
}
@@ -167,11 +159,7 @@ fn handle_start_message(
167159
rkyv::access::<protocol::ArchivedControlMessage, rkyv::rancor::Error>(received_data)?;
168160

169161
match archived_msg {
170-
protocol::ArchivedControlMessage::Start {
171-
buffer_size,
172-
timestamp_type,
173-
..
174-
} => {
162+
protocol::ArchivedControlMessage::Start { buffer_size, args } => {
175163
if ctx.another_started {
176164
let error_msg = "another client already started".to_string();
177165
let nack_msg = ControlMessage::Nack { error: &error_msg };
@@ -181,12 +169,31 @@ fn handle_start_message(
181169
}
182170

183171
let buffer_size = buffer_size.to_native() as usize;
184-
let timestamp_type_native = match *timestamp_type {
185-
protocol::ArchivedTimestampType::Monotonic => TimestampType::Monotonic,
186-
protocol::ArchivedTimestampType::Boottime => TimestampType::Boottime,
187-
protocol::ArchivedTimestampType::Realtime => TimestampType::Realtime,
172+
173+
let deserialized_args = match args {
174+
protocol::ArchivedArgs::Tracing(tracing_args) => {
175+
let timestamp_type_native = match tracing_args.timestamp_type {
176+
protocol::ArchivedTimestampType::Monotonic => TimestampType::Monotonic,
177+
protocol::ArchivedTimestampType::Boottime => TimestampType::Boottime,
178+
protocol::ArchivedTimestampType::Realtime => TimestampType::Realtime,
179+
};
180+
let log_filter = tracing_args.log_filter.as_str().to_string();
181+
182+
// Validate the log filter before proceeding
183+
if let Err(e) = tracing_subscriber::EnvFilter::try_new(&log_filter) {
184+
let error_msg = format!("Invalid log filter '{}': {}", log_filter, e);
185+
let nack_msg = ControlMessage::Nack { error: &error_msg };
186+
send_message(stream, &nack_msg)?;
187+
debug!("sent nack - {}", error_msg);
188+
return Ok(None);
189+
}
190+
191+
Some(protocol::Args::Tracing(protocol::TracingArgs {
192+
log_filter,
193+
timestamp_type: timestamp_type_native,
194+
}))
195+
}
188196
};
189-
let clock_id = timestamp_type_to_clock_id(timestamp_type_native);
190197

191198
let mut memory_fd: Option<OwnedFd> = None;
192199
let mut notification_fd: Option<OwnedFd> = None;
@@ -239,7 +246,7 @@ fn handle_start_message(
239246
}
240247
}
241248

242-
let producer_state = Arc::new(ProducerState::new(producer, clock_id));
249+
let producer_state = Arc::new(ProducerState::new(producer, deserialized_args));
243250

244251
debug!(buffer_size, client_id, "producer initialized");
245252

@@ -487,7 +494,7 @@ mod tests {
487494
use super::*;
488495
use crate::Consumer;
489496
use arc_swap::ArcSwapOption;
490-
use protocol::{ControlMessage, LogLevel};
497+
use protocol::ControlMessage;
491498
use rstest::{fixture, rstest};
492499
use std::os::unix::net::{UnixListener, UnixStream};
493500
use tempfile::TempDir;
@@ -519,11 +526,13 @@ mod tests {
519526
}
520527
}
521528

522-
fn send_start_message(stream: &UnixStream, buffer_size: u64, log_level: LogLevel) {
529+
fn send_start_message(stream: &UnixStream, buffer_size: u64, log_filter: &str) {
523530
let start_msg = ControlMessage::Start {
524531
buffer_size,
525-
log_level,
526-
timestamp_type: protocol::TimestampType::Monotonic,
532+
args: protocol::Args::Tracing(protocol::TracingArgs {
533+
log_filter: log_filter.to_string(),
534+
timestamp_type: protocol::TimestampType::Monotonic,
535+
}),
527536
};
528537

529538
let serialized_len = protocol::compute_length(&start_msg).unwrap();
@@ -552,8 +561,10 @@ mod tests {
552561
) {
553562
let start_msg = ControlMessage::Start {
554563
buffer_size: consumer.data_size() as u64,
555-
log_level: LogLevel::Debug,
556-
timestamp_type,
564+
args: protocol::Args::Tracing(protocol::TracingArgs {
565+
log_filter: "debug".to_string(),
566+
timestamp_type,
567+
}),
557568
};
558569

559570
let serialized_len = protocol::compute_length(&start_msg).unwrap();
@@ -638,7 +649,7 @@ mod tests {
638649
another_started: false,
639650
};
640651

641-
send_start_message(&client, 1024, LogLevel::Debug);
652+
send_start_message(&client, 1024, "debug");
642653
std::thread::sleep(std::time::Duration::from_millis(10));
643654

644655
let result = client_state.handle_message(&mut ctx);

‎agent/src/client.rs‎

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
use crate::Consumer;
22
use nix::sys::epoll::{Epoll, EpollCreateFlags, EpollEvent, EpollFlags, EpollTimeout};
33
use nix::sys::socket::{recv, sendmsg, ControlMessage, MsgFlags};
4-
use protocol::{ControlMessage as ProtocolControlMessage, LogLevel, TimestampType};
4+
use protocol::{Args, ControlMessage as ProtocolControlMessage, TimestampType, TracingArgs};
55
use std::os::fd::AsRawFd;
66
use std::os::unix::net::UnixStream;
77
use tracing::debug;
@@ -72,15 +72,15 @@ impl AgentClient {
7272
}
7373

7474
/// Connect to the agent and share the consumer's memory.
75-
pub fn start(&mut self, consumer: &Consumer, log_level: LogLevel) -> Result<()> {
76-
self.start_with_timestamp(consumer, log_level, TimestampType::Monotonic)
75+
pub fn start(&mut self, consumer: &Consumer, log_filter: String) -> Result<()> {
76+
self.start_with_timestamp(consumer, log_filter, TimestampType::Monotonic)
7777
}
7878

7979
/// Connect to the agent and share the consumer's memory with custom timestamp type.
8080
pub fn start_with_timestamp(
8181
&mut self,
8282
consumer: &Consumer,
83-
log_level: LogLevel,
83+
log_filter: String,
8484
timestamp_type: TimestampType,
8585
) -> Result<()> {
8686
let stream = UnixStream::connect(&self.socket_path)?;
@@ -89,8 +89,10 @@ impl AgentClient {
8989

9090
let start_msg = ProtocolControlMessage::Start {
9191
buffer_size: consumer.data_size() as u64,
92-
log_level,
93-
timestamp_type,
92+
args: Args::Tracing(TracingArgs {
93+
log_filter,
94+
timestamp_type,
95+
}),
9496
};
9597

9698
let serialized_len = protocol::compute_length(&start_msg)?;

‎agent/src/lib.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,8 @@ pub enum AgentError {
5959
Archive(#[from] rkyv::rancor::Error),
6060
#[error("Agent not enabled")]
6161
NotEnabled,
62+
#[error("{0}")]
63+
Other(String),
6264
}
6365

6466
pub type Result<T> = std::result::Result<T, AgentError>;

‎config.toml‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,4 +16,5 @@ filter_process = ["continuous", "manytrace"]
1616
socket = "/tmp/user1.sock"
1717

1818
[[user]]
19-
socket = "/tmp/user2.sock"
19+
socket = "/tmp/user2.sock"
20+
log_filter = "INFO"

‎manytrace/src/bin/manytrace.rs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ use clap::Parser;
33
use eyre::{Context, Result};
44
use manytrace::config::Config;
55
use manytrace::converter::PerfettoConverter;
6-
use protocol::{Event, LogLevel};
6+
use protocol::Event;
77
use std::cell::RefCell;
88
use std::fs::File;
99
use std::io::BufWriter;
@@ -79,7 +79,7 @@ fn main() -> Result<()> {
7979
if let Some(ref consumer) = user_consumer {
8080
for user_config in &config.user {
8181
let mut client = AgentClient::new(user_config.socket.clone());
82-
match client.start(consumer, LogLevel::Debug) {
82+
match client.start(consumer, user_config.log_filter.clone()) {
8383
Ok(_) => {
8484
tracing::info!(socket = %user_config.socket, "connected to user agent");
8585
user_clients.push(client);

‎manytrace/src/config.rs‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,12 @@ pub struct GlobalConfig {
2222
#[derive(Debug, Serialize, Deserialize)]
2323
pub struct UserConfig {
2424
pub socket: String,
25+
#[serde(default = "default_log_filter")]
26+
pub log_filter: String,
27+
}
28+
29+
fn default_log_filter() -> String {
30+
"DEBUG".to_string()
2531
}
2632

2733
fn default_buffer_size() -> usize {

0 commit comments

Comments
 (0)