datadog_agent_remote_config/
worker.rs

1use 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
20/// How often the worker retries until its first successful poll, because some products may not start correctly without
21/// a first state.
22const FIRST_POLL_RETRY: Duration = Duration::from_secs(1);
23
24/// Drives polling and delivery for a [`RemoteConfigurationClient`](crate::RemoteConfigurationClient).
25///
26/// A supervisor can own and restart the worker, or the caller can drive it directly with [`run`](Self::run).
27///
28/// The worker does not exit on connection failures, on the Agent having remote configuration disabled, on the Agent
29/// reporting its configuration expired, or on a panic in a subscriber's decoder. It keeps polling through all of them,
30/// so it exits only on a failure in the client itself.
31///
32/// A restart keeps every subscription and each product's last accepted snapshot, and discards the protocol state, so
33/// the restarted worker fetches and decodes everything again. Subscribers may therefore see a snapshot equal to the one
34/// they already hold.
35///
36/// Failed polls, responses the client cannot apply, and rejected configurations are counted in the client's metrics.
37/// Subscribers see only their own decoding errors, so they need the metrics or worker logs to see other failures.
38pub struct RemoteConfigurationWorker {
39    /// The subscriptions and client ID, shared with every client handle and kept across restarts.
40    pub(crate) shared: Arc<Shared>,
41
42    /// The Agent connection, kept across restarts and locked by the one running poll loop.
43    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    /// Runs the worker without a supervisor.
67    ///
68    /// This polls until the returned future is dropped.
69    ///
70    /// # Errors
71    ///
72    /// Returns an error only on a failure in the client itself. Connection failures, remote configuration being
73    /// disabled on the Agent, and rejected or panicking decoders are handled without returning.
74    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
95/// Polls until `shutdown` resolves, starting from fresh protocol state.
96async 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    // The reference point until the first success, so the gauge measures how long the worker has gone without a
108    // successful poll from the moment it started.
109    let mut last_success = tokio::time::Instant::now();
110
111    loop {
112        // This poll includes every subscription made so far, so a wake-up already pending for one of them is spent.
113        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                    // `GenericError::to_string` omits the cause chain, so two errors with the same outer message but
142                    // different causes would otherwise look identical here.
143                    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        // An `Unimplemented` answer is definitive, so it is not staleness; the polls counter still shows that Remote
165        // Configuration is off.
166        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/// A failed poll, as compared against the one before it to decide whether it is a repeat.
183#[derive(PartialEq)]
184enum Failure {
185    /// An RPC error with any message: a changed message during one outage is still the same outage.
186    Rpc,
187
188    /// An invalid response, with its full error chain.
189    InvalidResponse(String),
190
191    Unimplemented,
192}
193
194/// Decides how long to wait before the next poll, and whether a failure repeats the one before it.
195struct Schedule {
196    poll_interval: Duration,
197    max_backoff: Duration,
198    backoff: ExponentialBackoff,
199    succeeded_once: bool,
200    failures: u32,
201
202    /// The most recent failure since the last success. A failure identical to it logs at debug; any other logs at its
203    /// normal level.
204    last_failure: Option<Failure>,
205
206    /// Whether the Agent answered `Unimplemented` since the last success, so that the next success says polling
207    /// resumed rather than recovered.
208    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    /// Records `failure` as the most recent one and returns whether it is identical to the one it replaces.
250    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    /// Remote Configuration is disabled on the Agent, which is expected to last, so the worker keeps checking slowly.
257    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}