@@ -11,6 +11,10 @@ use std::{
1111} ;
1212
1313use anyhow:: { anyhow, Context } ;
14+ use datadog_sidecar:: service:: {
15+ telemetry:: { InProcessTelemetryClient , InProcessTelemetryClientFactory } ,
16+ InstanceId ,
17+ } ;
1418use futures:: stream:: Stream ;
1519use libddwaf:: { object:: WafObjectType , RunnableContext } ;
1620use log:: { debug, error, info, warning as warn} ;
@@ -92,19 +96,24 @@ pub struct Client {
9296 pub id : u64 ,
9397 service_manager : & ' static ServiceManager ,
9498 service : Option < TrackedService > ,
95- sidecar_settings : Option < protocol:: SidecarSettings > ,
99+ telemetry_client_factory : InProcessTelemetryClientFactory ,
100+ telemetry_client : Option < InProcessTelemetryClient > ,
96101 metrics_last_registered : Cell < Option < Instant > > ,
97102}
98103
99104static CLIENT_SERIAL : AtomicU64 = AtomicU64 :: new ( 1 ) ;
100105impl Client {
101- pub fn new ( service_manager : & ' static ServiceManager ) -> Self {
106+ pub fn new (
107+ service_manager : & ' static ServiceManager ,
108+ telemetry_client_factory : InProcessTelemetryClientFactory ,
109+ ) -> Self {
102110 Self {
103111 id : CLIENT_SERIAL . fetch_add ( 1 , atomic:: Ordering :: Relaxed ) ,
104112 service_manager,
105113 service : None ,
106- sidecar_settings : None ,
107- metrics_last_registered : Default :: default ( ) ,
114+ telemetry_client_factory,
115+ telemetry_client : None ,
116+ metrics_last_registered : Cell :: new ( None ) ,
108117 }
109118 }
110119
@@ -273,21 +282,26 @@ fn handle_client_init(
273282 args. telemetry_settings ,
274283 ) ;
275284
276- client. sidecar_settings = Some ( args. sidecar_settings . clone ( ) ) ;
285+ let telemetry_client = client. telemetry_client_factory . create_client (
286+ InstanceId :: new (
287+ args. sidecar_settings . session_id . clone ( ) ,
288+ args. sidecar_settings . runtime_id . clone ( ) ,
289+ ) ,
290+ telemetry_settings. service_name . clone ( ) ,
291+ telemetry_settings. env_name . clone ( ) ,
292+ ) ;
277293
278- update_error_telemetry_context ( args . sidecar_settings . clone ( ) , telemetry_settings . clone ( ) ) ;
294+ update_error_telemetry_context ( telemetry_client . clone ( ) ) ;
279295
280- let last_registration_time = & client. metrics_last_registered ;
281- let mut tel_metric_submitter = TelemetrySidecarMetricSubmitter :: create (
282- & args. sidecar_settings ,
283- & telemetry_settings,
284- last_registration_time,
285- ) ;
296+ let mut tel_metric_submitter =
297+ TelemetrySidecarMetricSubmitter :: create ( & telemetry_client, & client. metrics_last_registered ) ;
286298
287299 let service = client
288300 . service_manager
289301 . get_service ( & sd, tel_metric_submitter. as_mut ( ) ) ;
290302
303+ drop ( tel_metric_submitter) ;
304+
291305 let mut cir = ClientInitResp {
292306 version : protocol:: VERSION_FOR_PROTO ,
293307 client_id : client. id ,
@@ -307,10 +321,12 @@ fn handle_client_init(
307321 cir. helper_runtime = Some ( "rust" . to_string ( ) ) ;
308322
309323 client. service = Some ( TrackedService :: new ( service) ) ;
324+ client. telemetry_client = Some ( telemetry_client) ;
310325 cir. status = "ok" . to_string ( ) ;
311326 Ok ( CommandResponse :: ClientInit ( cir) )
312327 }
313328 Err ( err) => {
329+ clear_error_telemetry_context ( ) ;
314330 error ! ( "client init handling error: {:?}" , err) ;
315331 Err ( err) . context ( "client init handling error" )
316332 }
@@ -352,37 +368,40 @@ fn handle_config_sync(client: &mut Client, args: protocol::ConfigSyncArgs) {
352368 submit_service_telemetry ( client, service) ;
353369 }
354370
355- // potentially we have new telemetry settings, update the error telemetry context
356- if let Some ( ref sidecar_settings) = client. sidecar_settings {
357- update_error_telemetry_context ( sidecar_settings. clone ( ) , telemetry_settings. clone ( ) ) ;
358- } else {
371+ let Some ( new_telemetry) = client. telemetry_client . as_ref ( ) . map ( |telemetry| {
372+ telemetry. with_new_service_env (
373+ telemetry_settings. service_name . clone ( ) ,
374+ telemetry_settings. env_name . clone ( ) ,
375+ )
376+ } ) else {
377+ error ! ( "Cannot update telemetry client: client_init telemetry is missing" ) ;
359378 clear_error_telemetry_context ( ) ;
360- }
379+ return ;
380+ } ;
381+ client. metrics_last_registered . set ( None ) ;
382+ update_error_telemetry_context ( new_telemetry. clone ( ) ) ;
361383
362384 // ... and create a new telemetry metrics submitter because creating a new
363385 // service generates telemetry (more likely we're not creating a new
364386 // service though, we're just fetching the existing one)
365- let mut tel_metric_submitter = match client. sidecar_settings {
366- Some ( ref sidecar_settings) => TelemetrySidecarMetricSubmitter :: create (
367- sidecar_settings,
368- & telemetry_settings,
369- & client. metrics_last_registered ,
370- ) ,
371- None => {
372- // this should have been set in client_init
373- error ! ( "Cannot submit telemetry metrics: sidecar_settings unexpectadly not set" ) ;
374- TelemetrySidecarMetricSubmitter :: noop ( )
375- }
376- } ;
387+ let mut tel_metric_submitter =
388+ TelemetrySidecarMetricSubmitter :: create ( & new_telemetry, & client. metrics_last_registered ) ;
377389
378- match client
390+ let new_service = client
379391 . service_manager
380- . get_service ( & new_disc, & mut * tel_metric_submitter)
381- {
392+ . get_service ( & new_disc, & mut * tel_metric_submitter) ;
393+
394+ drop ( tel_metric_submitter) ;
395+
396+ match new_service {
382397 Ok ( new_service) => {
383398 client. service = Some ( TrackedService :: new ( new_service) ) ;
399+ client. telemetry_client = Some ( new_telemetry) ;
384400 }
385401 Err ( e) => {
402+ if let Some ( telemetry) = & client. telemetry_client {
403+ update_error_telemetry_context ( telemetry. clone ( ) ) ;
404+ }
386405 error ! ( "Failed to get service with new RC path: {}; will continue running with old service!" , e) ;
387406 }
388407 }
@@ -582,42 +601,26 @@ async fn run_request(
582601}
583602
584603fn submit_service_telemetry ( client : & Client , service : & Service ) {
585- if let ( Some ( sidecar_settings) , telemetry_settings) =
586- ( & client. sidecar_settings , & service. telemetry_settings ( ) )
587- {
604+ if let Some ( telemetry) = & client. telemetry_client {
588605 debug ! ( "Submitting service telemetry to sidecar" ) ;
589- let mut submitter =
590- TelemetrySidecarLogSubmitter :: create ( sidecar_settings, telemetry_settings) ;
606+ let mut submitter = TelemetrySidecarLogSubmitter :: create ( telemetry) ;
591607 service. generate_telemetry_logs ( & mut * submitter) ;
592608
593- let mut submitter = TelemetrySidecarMetricSubmitter :: create (
594- sidecar_settings,
595- telemetry_settings,
596- & client. metrics_last_registered ,
597- ) ;
609+ let mut submitter =
610+ TelemetrySidecarMetricSubmitter :: create ( telemetry, & client. metrics_last_registered ) ;
598611 service. generate_telemetry_metrics ( & mut * submitter) ;
599612 } else {
600- debug ! (
601- "Cannot submit service telemetry: sidecar_settings={:?}, telemetry_settings={:?}" ,
602- client. sidecar_settings,
603- service. telemetry_settings( )
604- ) ;
613+ debug ! ( "Cannot submit service telemetry: telemetry client is not bound" ) ;
605614 }
606615}
607616
608617fn submit_context_telemetry_metrics ( client : & Client , req_ctx : & mut ReqContext ) {
609- let Some ( ref sidecar_settings ) = client. sidecar_settings else {
610- warn ! ( "Cannot submit context telemetry metrics: sidecar_settings not set " ) ;
618+ let Some ( telemetry ) = & client. telemetry_client else {
619+ warn ! ( "Cannot submit context telemetry metrics: telemetry client is not bound " ) ;
611620 return ;
612621 } ;
613- let service = client. get_service ( ) ;
614- let telemetry_settings = service. telemetry_settings ( ) ;
615-
616- let mut tel_metric_submitter = TelemetrySidecarMetricSubmitter :: create (
617- sidecar_settings,
618- telemetry_settings,
619- & client. metrics_last_registered ,
620- ) ;
622+ let mut tel_metric_submitter =
623+ TelemetrySidecarMetricSubmitter :: create ( telemetry, & client. metrics_last_registered ) ;
621624
622625 let waf_metrics = req_ctx. take_waf_metrics ( ) ;
623626 waf_metrics. generate_telemetry_metrics ( & mut * tel_metric_submitter) ;
0 commit comments