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

Commit a21006d

Browse files
committed
tracing-manytrace: add thread and process name events
Introduce ThreadName and ProcessName event types that are emitted once per thread/process when a new client connects. This allows consumers to identify threads and processes by name rather than just by ID.
1 parent 507b183 commit a21006d

6 files changed

Lines changed: 268 additions & 62 deletions

File tree

‎agent/src/agent.rs‎

Lines changed: 29 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -7,9 +7,24 @@ use std::thread::{self, JoinHandle};
77
use crate::epoll_thread::epoll_listener_thread;
88
use crate::{Producer, Result};
99

10+
/// Producer state containing producer and client_id.
11+
pub(crate) struct ProducerState {
12+
producer: Arc<Producer>,
13+
client_id: u64,
14+
}
15+
16+
impl ProducerState {
17+
pub(crate) fn new(producer: Arc<Producer>, client_id: u64) -> Self {
18+
Self {
19+
producer,
20+
client_id,
21+
}
22+
}
23+
}
24+
1025
/// Server that listens for client connections and receives events.
1126
pub struct Agent {
12-
producer: Arc<ArcSwapOption<Producer>>,
27+
producer_state: Arc<ArcSwapOption<ProducerState>>,
1328
socket_thread: Option<JoinHandle<Result<()>>>,
1429
socket_path: String,
1530
shutdown: Arc<AtomicBool>,
@@ -26,8 +41,8 @@ impl Agent {
2641

2742
/// Create a new agent with custom client keepalive timeout.
2843
pub fn with_timeout(socket_path: String, timeout_ms: u16) -> Result<Self> {
29-
let producer = Arc::new(ArcSwapOption::empty());
30-
let producer_clone = producer.clone();
44+
let producer_state = Arc::new(ArcSwapOption::empty());
45+
let producer_state_clone = producer_state.clone();
3146
let socket_path_clone = socket_path.clone();
3247
let shutdown = Arc::new(AtomicBool::new(false));
3348
let shutdown_clone = shutdown.clone();
@@ -37,14 +52,14 @@ impl Agent {
3752
.spawn(move || {
3853
epoll_listener_thread(
3954
socket_path_clone,
40-
producer_clone,
55+
producer_state_clone,
4156
shutdown_clone,
4257
timeout_ms,
4358
)
4459
})?;
4560

4661
Ok(Agent {
47-
producer,
62+
producer_state,
4863
socket_thread: Some(socket_thread),
4964
socket_path,
5065
shutdown,
@@ -53,18 +68,23 @@ impl Agent {
5368

5469
/// Check if a client is currently connected.
5570
pub fn enabled(&self) -> bool {
56-
self.producer.load().is_some()
71+
self.producer_state.load().is_some()
5772
}
5873

5974
/// Submit an event to connected clients.
6075
pub fn submit(&self, event: &Event) -> Result<()> {
61-
let producer = self.producer.load();
62-
match producer.as_ref() {
63-
Some(producer) => producer.submit(event),
76+
let producer_state = self.producer_state.load();
77+
match producer_state.as_ref() {
78+
Some(state) => state.producer.submit(event),
6479
None => Err(crate::AgentError::NotEnabled),
6580
}
6681
}
6782

83+
/// Get the current client_id if available.
84+
pub fn client_id(&self) -> Option<u64> {
85+
self.producer_state.load().as_ref().map(|s| s.client_id)
86+
}
87+
6888
/// Gracefully shut down the agent and wait for thread termination.
6989
pub fn wait_terminated(&mut self) -> Result<()> {
7090
self.shutdown.store(true, Ordering::Relaxed);

‎agent/src/agent_state.rs‎

Lines changed: 19 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ use std::os::unix::net::UnixStream;
88
use std::sync::Arc;
99
use tracing::{debug, warn};
1010

11+
use crate::agent::ProducerState;
1112
use crate::{AgentError, Producer, Result};
1213

1314
pub(crate) struct AgentContext {
@@ -62,7 +63,7 @@ impl AgentClientState {
6263

6364
pub(crate) enum MessageResult {
6465
Continue,
65-
Started(Arc<Producer>),
66+
Started(Arc<ProducerState>),
6667
Disconnect,
6768
#[allow(dead_code)]
6869
Error(AgentError),
@@ -72,12 +73,12 @@ fn handle_pending_state(
7273
client_state: &mut AgentClientState,
7374
ctx: &mut AgentContext,
7475
) -> MessageResult {
75-
match handle_start_message(&client_state.stream, ctx) {
76-
Ok(Some(producer)) => {
76+
match handle_start_message(&client_state.stream, ctx, client_state.client_id) {
77+
Ok(Some(producer_state_result)) => {
7778
client_state.timestamp_ns = get_timestamp_ns();
7879
client_state.state = State::Started;
7980
debug!(client_id = client_state.client_id, "client started");
80-
MessageResult::Started(producer)
81+
MessageResult::Started(producer_state_result)
8182
}
8283
Ok(None) => MessageResult::Continue,
8384
Err(e) => {
@@ -132,7 +133,8 @@ fn handle_started_state(client_state: &mut AgentClientState) -> MessageResult {
132133
fn handle_start_message(
133134
stream: &UnixStream,
134135
ctx: &mut AgentContext,
135-
) -> Result<Option<Arc<Producer>>> {
136+
client_id: u64,
137+
) -> Result<Option<Arc<ProducerState>>> {
136138
let mut cmsg_buffer = nix::cmsg_space!([RawFd; 2]);
137139
let mut msg_buf = [0u8; 1024];
138140
let mut iov = [std::io::IoSliceMut::new(&mut msg_buf)];
@@ -202,14 +204,15 @@ fn handle_start_message(
202204
)?;
203205

204206
let producer = Arc::new(Producer::from_inner(new_producer));
207+
let producer_state = Arc::new(ProducerState::new(producer, client_id));
205208

206-
debug!(buffer_size, "producer initialized");
209+
debug!(buffer_size, client_id, "producer initialized");
207210

208211
let ack_msg = ControlMessage::Ack;
209212
send_message(stream, &ack_msg)?;
210213
debug!("sent ack - producer initialized");
211214

212-
Ok(Some(producer))
215+
Ok(Some(producer_state))
213216
}
214217
_ => {
215218
debug!("expected Start message as first message");
@@ -307,16 +310,16 @@ impl AgentState {
307310
pub(crate) fn handle_event(
308311
&mut self,
309312
event_data: u64,
310-
producer: &Arc<ArcSwapOption<Producer>>,
313+
producer_state: &Arc<ArcSwapOption<ProducerState>>,
311314
) -> Option<EpollAction> {
312315
if let Some(mut client) = self.pending_clients.remove(&event_data) {
313316
match client.handle_message(&mut self.ctx) {
314317
MessageResult::Continue => {
315318
self.pending_clients.insert(event_data, client);
316319
None
317320
}
318-
MessageResult::Started(new_producer) => {
319-
producer.store(Some(new_producer));
321+
MessageResult::Started(new_producer_state) => {
322+
producer_state.store(Some(new_producer_state));
320323
self.ctx.another_started = true;
321324
self.started_client = Some((event_data, client));
322325
None
@@ -344,7 +347,7 @@ impl AgentState {
344347
}
345348
MessageResult::Disconnect | MessageResult::Error(_) => {
346349
let fd = client.into_stream();
347-
producer.store(None);
350+
producer_state.store(None);
348351
self.ctx.another_started = false;
349352
Some(EpollAction::Delete { fd })
350353
}
@@ -360,7 +363,10 @@ impl AgentState {
360363
}
361364
}
362365

363-
pub(crate) fn timer(&mut self, producer: &Arc<ArcSwapOption<Producer>>) -> Vec<EpollAction> {
366+
pub(crate) fn timer(
367+
&mut self,
368+
producer_state: &Arc<ArcSwapOption<ProducerState>>,
369+
) -> Vec<EpollAction> {
364370
debug!(
365371
pending_clients = self.pending_clients.len(),
366372
has_started_client = self.started_client.is_some(),
@@ -421,7 +427,7 @@ impl AgentState {
421427
});
422428
}
423429
self.started_client = None;
424-
producer.store(None);
430+
producer_state.store(None);
425431
self.ctx.another_started = false;
426432
}
427433
}

‎agent/src/epoll_thread.rs‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,15 @@ use std::sync::atomic::{AtomicBool, Ordering};
55
use std::sync::Arc;
66
use tracing::{debug, warn};
77

8+
use crate::agent::ProducerState;
89
use crate::agent_state::{AgentState, EpollAction};
9-
use crate::{Producer, Result};
10+
use crate::Result;
1011

1112
const LISTENER_TOKEN: u64 = 0;
1213

1314
pub fn epoll_listener_thread(
1415
socket_path: String,
15-
producer: Arc<ArcSwapOption<Producer>>,
16+
producer_state: Arc<ArcSwapOption<ProducerState>>,
1617
shutdown: Arc<AtomicBool>,
1718
timeout_ms: u16,
1819
) -> Result<()> {
@@ -59,7 +60,7 @@ pub fn epoll_listener_thread(
5960
}
6061
},
6162
client_id if client_id != LISTENER_TOKEN => {
62-
if let Some(action) = agent_state.handle_event(client_id, &producer) {
63+
if let Some(action) = agent_state.handle_event(client_id, &producer_state) {
6364
match action {
6465
EpollAction::Add { fd, token } => {
6566
epoll.add(fd, EpollEvent::new(EpollFlags::EPOLLIN, token))?;
@@ -73,7 +74,7 @@ pub fn epoll_listener_thread(
7374
_ => {}
7475
}
7576
}
76-
let timeout_actions = agent_state.timer(&producer);
77+
let timeout_actions = agent_state.timer(&producer_state);
7778
for action in timeout_actions {
7879
match action {
7980
EpollAction::Add { .. } => {}

‎protocol/src/lib.rs‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -79,6 +79,21 @@ pub struct Instant<'a> {
7979
pub labels: Labels<'a>,
8080
}
8181

82+
#[derive(Archive, Serialize, Deserialize)]
83+
pub struct ThreadName<'a> {
84+
#[rkyv(with = InlineAsBox)]
85+
pub name: &'a str,
86+
pub tid: i32,
87+
pub pid: i32,
88+
}
89+
90+
#[derive(Archive, Serialize, Deserialize)]
91+
pub struct ProcessName<'a> {
92+
#[rkyv(with = InlineAsBox)]
93+
pub name: &'a str,
94+
pub pid: i32,
95+
}
96+
8297
#[derive(Archive, Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq)]
8398
#[rkyv(compare(PartialEq), derive(Debug))]
8499
pub enum LogLevel {
@@ -94,6 +109,8 @@ pub enum Event<'a> {
94109
Counter(Counter<'a>),
95110
Span(Span<'a>),
96111
Instant(Instant<'a>),
112+
ThreadName(ThreadName<'a>),
113+
ProcessName(ProcessName<'a>),
97114
}
98115

99116
#[derive(Archive, Serialize, Deserialize, Debug, Clone, PartialEq, Eq)]
@@ -404,6 +421,43 @@ mod tests {
404421
}
405422
}
406423

424+
#[test]
425+
fn test_thread_process_name_serialization() {
426+
let thread_name = ThreadName {
427+
name: "worker-thread-1",
428+
tid: 12345,
429+
pid: 67890,
430+
};
431+
432+
let process_name = ProcessName {
433+
name: "my-process",
434+
pid: 67890,
435+
};
436+
437+
let thread_event = Event::ThreadName(thread_name);
438+
let process_event = Event::ProcessName(process_name);
439+
440+
for event in [thread_event, process_event] {
441+
let buf =
442+
to_bytes_in::<_, Error>(&event, Vec::new()).expect("event serialization failed");
443+
let archived = rkyv::access::<ArchivedEvent, rkyv::rancor::Error>(&buf)
444+
.expect("failed to access archived event");
445+
446+
match (&event, archived) {
447+
(Event::ThreadName(thread), ArchivedEvent::ThreadName(arch_thread)) => {
448+
assert_eq!(thread.name.as_bytes(), arch_thread.name.as_bytes());
449+
assert_eq!(thread.tid, arch_thread.tid.to_native());
450+
assert_eq!(thread.pid, arch_thread.pid.to_native());
451+
}
452+
(Event::ProcessName(process), ArchivedEvent::ProcessName(arch_process)) => {
453+
assert_eq!(process.name.as_bytes(), arch_process.name.as_bytes());
454+
assert_eq!(process.pid, arch_process.pid.to_native());
455+
}
456+
_ => panic!("mismatched event variants"),
457+
}
458+
}
459+
}
460+
407461
#[test]
408462
fn test_control_message_serialization() {
409463
let messages = [

0 commit comments

Comments
 (0)