datadog_agent_remote_config/
worker.rs1use std::fmt;
2use std::pin::pin;
3use std::sync::Arc;
4use std::time::Duration;
5
6use async_trait::async_trait;
7use saluki_common::sync::shutdown::ShutdownHandle;
8use saluki_core::runtime::{InitializationError, Supervisable, SupervisorFuture};
9use saluki_error::{generic_error, GenericError};
10use saluki_io::net::util::retry::ExponentialBackoff;
11use tokio::sync::Mutex;
12use tracing::{debug, error, info, warn};
13
14use crate::metrics::{Metrics, PollOutcome};
15use crate::registry::Shared;
16use crate::repository::Repository;
17use crate::source::{FetchError, RcAgent};
18use crate::RcClientConfiguration;
19
20const FIRST_POLL_RETRY: Duration = Duration::from_secs(1);
23
24pub struct RemoteConfigurationWorker {
39 pub(crate) shared: Arc<Shared>,
41
42 agent: Arc<Mutex<Box<dyn RcAgent>>>,
44
45 config: RcClientConfiguration,
46}
47
48impl fmt::Debug for RemoteConfigurationWorker {
49 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
50 f.debug_struct("RemoteConfigurationWorker")
51 .field("client_id", &self.shared.client_id)
52 .field("config", &self.config)
53 .finish_non_exhaustive()
54 }
55}
56
57impl RemoteConfigurationWorker {
58 pub(crate) fn new(shared: Arc<Shared>, agent: Box<dyn RcAgent>, config: RcClientConfiguration) -> Self {
59 Self {
60 shared,
61 agent: Arc::new(Mutex::new(agent)),
62 config,
63 }
64 }
65
66 pub async fn run(self) -> Result<(), GenericError> {
75 poll_loop(self.shared, self.agent, self.config, ShutdownHandle::noop()).await
76 }
77}
78
79#[async_trait]
80impl Supervisable for RemoteConfigurationWorker {
81 fn name(&self) -> &str {
82 "remote-config"
83 }
84
85 async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
86 Ok(Box::pin(poll_loop(
87 Arc::clone(&self.shared),
88 Arc::clone(&self.agent),
89 self.config.clone(),
90 process_shutdown,
91 )))
92 }
93}
94
95async fn poll_loop(
97 shared: Arc<Shared>, agent: Arc<Mutex<Box<dyn RcAgent>>>, config: RcClientConfiguration, shutdown: ShutdownHandle,
98) -> Result<(), GenericError> {
99 let mut shutdown = pin!(shutdown);
100 let mut agent = tokio::select! {
101 _ = &mut shutdown => return Ok(()),
102 agent = agent.lock() => agent,
103 };
104 let mut repository = Repository::new();
105 let mut schedule = Schedule::new(&config);
106 let metrics = Metrics::new();
107 let mut last_success = tokio::time::Instant::now();
110
111 loop {
112 tokio::select! {
114 biased;
115 _ = shared.wake.notified() => {}
116 _ = std::future::ready(()) => {}
117 }
118
119 let live = shared.live_products();
120 repository.prune(&live);
121 let request = repository.request(&shared.client_id, &config, &live);
122 let result = tokio::select! {
123 _ = &mut shutdown => return Ok(()),
124 result = tokio::time::timeout(config.request_timeout, agent.get_configs(request)) => {
125 result.unwrap_or_else(|_| {
126 Err(FetchError::Rpc(generic_error!(
127 "The Agent did not answer within {:?}.",
128 config.request_timeout
129 )))
130 })
131 }
132 };
133
134 let (outcome, delay) = match result {
135 Ok(response) => match repository.apply(response, &live, &shared, &metrics) {
136 Ok(outcome) => {
137 repository.last_error = None;
138 (outcome, schedule.succeeded())
139 }
140 Err(e) => {
141 if schedule.repeats(Failure::InvalidResponse(format!("{e:#}"))) {
144 debug!(error = %e, "Discarded an invalid Remote Configuration response.");
145 } else {
146 error!(error = %e, "Discarded an invalid Remote Configuration response.");
147 }
148 repository.last_error = Some(e.to_string());
149 (PollOutcome::InvalidResponse, schedule.failed())
150 }
151 },
152 Err(FetchError::Unimplemented(e)) => (PollOutcome::Unimplemented, schedule.unimplemented(&e)),
153 Err(FetchError::Rpc(e)) => {
154 if schedule.repeats(Failure::Rpc) {
155 debug!(error = %e, "Failed to poll the Agent for Remote Configuration.");
156 } else {
157 warn!(error = %e, "Failed to poll the Agent for Remote Configuration.");
158 }
159 repository.last_error = Some(e.to_string());
160 (PollOutcome::RpcError, schedule.failed())
161 }
162 };
163 metrics.count_poll(outcome);
164 if matches!(
167 outcome,
168 PollOutcome::Ok | PollOutcome::Expired | PollOutcome::Unimplemented
169 ) {
170 last_success = tokio::time::Instant::now();
171 }
172 metrics.set_seconds_since_successful_poll(last_success.elapsed().as_secs());
173
174 tokio::select! {
175 _ = &mut shutdown => return Ok(()),
176 _ = tokio::time::sleep(delay) => {}
177 _ = shared.wake.notified() => {}
178 }
179 }
180}
181
182#[derive(PartialEq)]
184enum Failure {
185 Rpc,
187
188 InvalidResponse(String),
190
191 Unimplemented,
192}
193
194struct Schedule {
196 poll_interval: Duration,
197 max_backoff: Duration,
198 backoff: ExponentialBackoff,
199 succeeded_once: bool,
200 failures: u32,
201
202 last_failure: Option<Failure>,
205
206 unimplemented_seen: bool,
209}
210
211impl Schedule {
212 fn new(config: &RcClientConfiguration) -> Self {
213 Self {
214 poll_interval: config.poll_interval,
215 max_backoff: config.max_backoff,
216 backoff: ExponentialBackoff::with_jitter(config.poll_interval, config.max_backoff, 2.0),
217 succeeded_once: false,
218 failures: 0,
219 last_failure: None,
220 unimplemented_seen: false,
221 }
222 }
223
224 fn succeeded(&mut self) -> Duration {
225 if self.unimplemented_seen {
226 info!("Remote Configuration is enabled on the Agent; polling resumed.");
227 } else if self.failures > 0 {
228 info!(
229 failures = self.failures,
230 "Polling the Agent for Remote Configuration recovered."
231 );
232 }
233 self.succeeded_once = true;
234 self.failures = 0;
235 self.last_failure = None;
236 self.unimplemented_seen = false;
237 self.poll_interval
238 }
239
240 fn failed(&mut self) -> Duration {
241 self.failures = self.failures.saturating_add(1);
242 if self.succeeded_once {
243 self.backoff.get_backoff_duration(self.failures)
244 } else {
245 FIRST_POLL_RETRY
246 }
247 }
248
249 fn repeats(&mut self, failure: Failure) -> bool {
251 let repeated = self.last_failure.as_ref() == Some(&failure);
252 self.last_failure = Some(failure);
253 repeated
254 }
255
256 fn unimplemented(&mut self, error: &GenericError) -> Duration {
258 if !self.repeats(Failure::Unimplemented) {
259 info!(error = %error, "Remote Configuration is not enabled on the Agent; checking again periodically.");
260 }
261 self.unimplemented_seen = true;
262 self.max_backoff
263 }
264}