@@ -17,7 +17,7 @@ use a3s_observer_common::{
1717use anyhow:: Context as _;
1818use aya:: {
1919 maps:: { PerCpuArray , RingBuf } ,
20- programs:: { TracePoint , UProbe } ,
20+ programs:: { KProbe , TracePoint , UProbe } ,
2121 Ebpf ,
2222} ;
2323use std:: collections:: HashMap ;
@@ -63,7 +63,6 @@ async fn main() -> anyhow::Result<()> {
6363 let files = std:: env:: var_os ( "A3S_OBSERVER_FILES" ) . is_some ( ) ;
6464 let mut probes = vec ! [
6565 ( "exec" , "sys_enter_execve" ) ,
66- ( "proc_exit" , "sys_enter_exit_group" ) ,
6766 ( "tls_write" , "sys_enter_write" ) ,
6867 ( "tls_sendto" , "sys_enter_sendto" ) ,
6968 ( "connect" , "sys_enter_connect" ) ,
@@ -91,6 +90,14 @@ async fn main() -> anyhow::Result<()> {
9190 }
9291 }
9392 }
93+ // proc_exit is a do_exit kprobe (not a tracepoint): do_exit fires for EVERY task exit,
94+ // including signal-kills (crash / OOM) that sys_enter_exit_group never sees.
95+ match attach_kprobe ( & mut ebpf, "proc_exit" , "do_exit" ) {
96+ Ok ( ( ) ) => attached += 1 ,
97+ Err ( e) => {
98+ tracing:: warn!( error = %e, "proc_exit (do_exit kprobe) failed — exit signals unavailable" )
99+ }
100+ }
94101 if attached == 0 {
95102 anyhow:: bail!( "no eBPF probes could be attached" ) ;
96103 }
@@ -132,7 +139,7 @@ async fn main() -> anyhow::Result<()> {
132139 let classifier = SniClassifier ;
133140 let resolver = KubeResolver ; // cgroup→pod in k8s; falls back to comm on bare hosts
134141 // (pid,fd) -> peer, populated by connect, read by the TLS probe to fuse provider+peer.
135- let mut peers: HashMap < u64 , IpAddr > = HashMap :: new ( ) ;
142+ let mut peers: HashMap < u64 , ( IpAddr , u16 ) > = HashMap :: new ( ) ;
136143 // (pid,fd) -> (sni, provider, peer): recorded at ClientHello, read when the socket
137144 // closes (the in-kernel LlmEvent) to build the metric-bearing LlmCall.
138145 let mut llm_meta: HashMap < u64 , ( Option < String > , Option < Provider > , IpAddr ) > = HashMap :: new ( ) ;
@@ -232,6 +239,7 @@ async fn main() -> anyhow::Result<()> {
232239 event: AgentEvent :: ToolExec {
233240 pid: ev. pid,
234241 ppid: read_ppid( ev. pid) ,
242+ uid: ev. uid,
235243 argv: argv_of( & ev) ,
236244 cwd: read_cwd( ev. pid) ,
237245 } ,
@@ -246,6 +254,7 @@ async fn main() -> anyhow::Result<()> {
246254 event: AgentEvent :: ProcessExit {
247255 pid: ev. pid,
248256 exit_code: ev. exit_code,
257+ signal: ev. signal,
249258 } ,
250259 } ) ;
251260 }
@@ -257,14 +266,15 @@ async fn main() -> anyhow::Result<()> {
257266 if peers. len( ) > 8192 {
258267 peers. clear( ) ; // ponytail: crude cap; LRU if it ever matters
259268 }
260- peers. insert( sock_key( ev. pid, ev. fd) , peer) ;
269+ peers. insert( sock_key( ev. pid, ev. fd) , ( peer, ev . port ) ) ;
261270 emit( exporter. as_ref( ) , & mut stats, EnrichedEvent {
262271 identity: identity_for( & resolver, ev. pid, & ev. comm) ,
263272 provider: None ,
264273 event: AgentEvent :: Egress {
265274 pid: ev. pid,
266275 sni: None ,
267276 peer,
277+ port: ev. port,
268278 bytes: 0 ,
269279 } ,
270280 } ) ;
@@ -275,10 +285,10 @@ async fn main() -> anyhow::Result<()> {
275285 let len = ( ev. len as usize ) . min( ev. data. len( ) ) ;
276286 let sni = parse_sni( & ev. data[ ..len] ) ;
277287 // Correlated peer for this socket (the LLM endpoint).
278- let peer = peers
288+ let ( peer, port ) = peers
279289 . get( & sock_key( ev. pid, ev. fd) )
280290 . copied( )
281- . unwrap_or( UNKNOWN_PEER ) ;
291+ . unwrap_or( ( UNKNOWN_PEER , 0 ) ) ;
282292 let provider =
283293 sni. as_deref( ) . and_then( |h| classifier. classify( Some ( h) , peer) ) ;
284294 // Remember the call so the close event can build a metric-bearing
@@ -294,6 +304,7 @@ async fn main() -> anyhow::Result<()> {
294304 pid: ev. pid,
295305 sni,
296306 peer,
307+ port,
297308 bytes: ev. len as u64 ,
298309 } ,
299310 } ) ;
@@ -364,8 +375,25 @@ async fn main() -> anyhow::Result<()> {
364375 let len = ( ev. len as usize ) . min( ev. data. len( ) ) ;
365376 let content = String :: from_utf8_lossy( & ev. data[ ..len] ) . into_owned( ) ;
366377 if !content. is_empty( ) {
378+ let identity = identity_for( & resolver, ev. pid, & ev. comm) ;
379+ // Structured LLM telemetry (model/tokens) alongside the raw content.
380+ if let Some ( ( model, prompt_tokens, completion_tokens) ) =
381+ parse_llm_meta( & content)
382+ {
383+ emit( exporter. as_ref( ) , & mut stats, EnrichedEvent {
384+ identity: identity. clone( ) ,
385+ provider: None ,
386+ event: AgentEvent :: LlmApi {
387+ pid: ev. pid,
388+ is_request: ev. is_read == 0 ,
389+ model,
390+ prompt_tokens,
391+ completion_tokens,
392+ } ,
393+ } ) ;
394+ }
367395 emit( exporter. as_ref( ) , & mut stats, EnrichedEvent {
368- identity: identity_for ( & resolver , ev . pid , & ev . comm ) ,
396+ identity,
369397 provider: None ,
370398 event: AgentEvent :: SslContent {
371399 pid: ev. pid,
@@ -406,6 +434,17 @@ fn attach(ebpf: &mut Ebpf, prog: &str, category: &str, name: &str) -> anyhow::Re
406434 Ok ( ( ) )
407435}
408436
437+ fn attach_kprobe ( ebpf : & mut Ebpf , prog : & str , sym : & str ) -> anyhow:: Result < ( ) > {
438+ let p: & mut KProbe = ebpf
439+ . program_mut ( prog)
440+ . with_context ( || format ! ( "`{prog}` program not found" ) ) ?
441+ . try_into ( ) ?;
442+ p. load ( ) ?;
443+ p. attach ( sym, 0 )
444+ . with_context ( || format ! ( "attach kprobe {sym}" ) ) ?;
445+ Ok ( ( ) )
446+ }
447+
409448fn attach_uprobe ( ebpf : & mut Ebpf , prog : & str , sym : & str , target : & str ) -> anyhow:: Result < ( ) > {
410449 let p: & mut UProbe = ebpf
411450 . program_mut ( prog)
@@ -486,6 +525,7 @@ fn emit(exporter: &dyn Exporter, stats: &mut Stats, ev: EnrichedEvent) {
486525 AgentEvent :: FileDelete { .. } => stats. file += 1 ,
487526 AgentEvent :: LlmCall { .. } => stats. llm += 1 ,
488527 AgentEvent :: SslContent { .. } => stats. ssl += 1 ,
528+ AgentEvent :: LlmApi { .. } => stats. llm += 1 ,
489529 }
490530 exporter. export ( & ev) ;
491531}
@@ -560,9 +600,33 @@ fn parse_dns_qname(buf: &[u8]) -> Option<String> {
560600 ( !name. is_empty ( ) ) . then_some ( name)
561601}
562602
603+ /// Best-effort LLM-API fields from captured TLS plaintext: `"model"` from a request body, token
604+ /// `usage` from a response. None if absent (not an LLM call, or the bytes weren't captured).
605+ /// Consumes untrusted plaintext — every index is bounds-checked, must never panic.
606+ fn parse_llm_meta ( s : & str ) -> Option < ( Option < String > , Option < u32 > , Option < u32 > ) > {
607+ let model = json_str_after ( s, "\" model\" " ) ;
608+ let pt = json_num_after ( s, "\" prompt_tokens\" " ) ;
609+ let ct = json_num_after ( s, "\" completion_tokens\" " ) ;
610+ ( model. is_some ( ) || pt. is_some ( ) || ct. is_some ( ) ) . then_some ( ( model, pt, ct) )
611+ }
612+
613+ fn json_str_after ( s : & str , key : & str ) -> Option < String > {
614+ let rest = & s[ s. find ( key) ? + key. len ( ) ..] ; // find() ≤ len, +key.len() ≤ len → in-bounds
615+ let body = & rest[ rest. find ( '"' ) ? + 1 ..] ; // past the value's opening quote
616+ Some ( body[ ..body. find ( '"' ) ?] . to_owned ( ) )
617+ }
618+
619+ fn json_num_after ( s : & str , key : & str ) -> Option < u32 > {
620+ let rest = s[ s. find ( key) ? + key. len ( ) ..] . trim_start_matches ( [ ':' , ' ' , '\t' ] ) ;
621+ let end = rest
622+ . find ( |c : char | !c. is_ascii_digit ( ) )
623+ . unwrap_or ( rest. len ( ) ) ;
624+ rest. get ( ..end) ?. parse ( ) . ok ( )
625+ }
626+
563627#[ cfg( test) ]
564628mod tests {
565- use super :: { parse_dns_qname, parse_sni} ;
629+ use super :: { parse_dns_qname, parse_llm_meta , parse_sni} ;
566630
567631 #[ test]
568632 fn parses_sni_from_minimal_clienthello ( ) {
@@ -589,6 +653,16 @@ mod tests {
589653 assert_eq ! ( parse_sni( & [ ] ) , None ) ;
590654 }
591655
656+ #[ test]
657+ fn parse_llm_meta_extracts_model_and_tokens ( ) {
658+ let req = r#"POST /v1/chat/completions HTTP/1.1 ... {"model":"gpt-4o","messages":[{"role":"user","content":"hi"}]}"# ;
659+ assert_eq ! ( parse_llm_meta( req) . unwrap( ) . 0 . as_deref( ) , Some ( "gpt-4o" ) ) ;
660+ let resp = r#"{"id":"x","choices":[],"usage":{"prompt_tokens":12,"completion_tokens":34}}"# ;
661+ let ( _, pt, ct) = parse_llm_meta ( resp) . unwrap ( ) ;
662+ assert_eq ! ( ( pt, ct) , ( Some ( 12 ) , Some ( 34 ) ) ) ;
663+ assert ! ( parse_llm_meta( "just plaintext, no json fields here" ) . is_none( ) ) ;
664+ }
665+
592666 #[ test]
593667 fn parses_dns_query_name ( ) {
594668 let mut q = vec ! [ 0u8 ; 12 ] ; // header
0 commit comments