datadog_agent_remote_config/
lib.rs

1//! Provides a client for remote configuration.
2//!
3//! Configuration assigned to this client arrives through polling the Datadog Agent, which fetches it from the Datadog
4//! backend. A subscriber supplies a [`ProductDecoder`] that names a product and decodes its payloads, then receives
5//! typed snapshots through a [`Subscription`]. The crate keeps the client's identity, its protocol cursor (the last
6//! targets version), the cached paths it reports to the Agent, configuration path structure, and numeric apply states
7//! private.
8//!
9//! # Testing
10//!
11//! With the `test-util` feature enabled, [`TestPublisher`] creates a [`Subscription`] that a test publishes into by
12//! hand, either with finished snapshots and rejections or by running a decoder over payloads exactly as the client
13//! does. A component that takes a `Subscription` can therefore be tested without an Agent.
14//!
15//! # Trust
16//!
17//! The client performs no TUF signature verification. It trusts the Agent, reached over an authenticated local IPC
18//! channel, to have verified already. It does validate that each payload matches the length and SHA-256 hash published
19//! in the accompanying targets metadata, which guards against a bug in delivery.
20//!
21//! # Examples
22//!
23//! A decoder for a product whose payload is a message, a consumer of its subscription, and a test of both through
24//! [`TestPublisher`]:
25//!
26//! ```
27//! use datadog_agent_remote_config::{ConfigId, ProductDecoder, RemoteConfigurationClient, Subscription, TestPublisher};
28//!
29//! #[derive(Debug)]
30//! struct Example {
31//!     message: String,
32//! }
33//!
34//! /// Keeps the last valid message in ascending configuration ID order.
35//! #[derive(Default)]
36//! struct ExampleDecoder {
37//!     message: Option<String>,
38//! }
39//!
40//! impl ProductDecoder for ExampleDecoder {
41//!     const PRODUCT: &'static str = "EXAMPLE_PRODUCT";
42//!
43//!     type Snapshot = Example;
44//!     type Error = String;
45//!
46//!     fn decode(&mut self, _id: &ConfigId, payload: &[u8]) -> Result<(), Self::Error> {
47//!         let message = std::str::from_utf8(payload).map_err(|_| "Message is not UTF-8.".to_owned())?;
48//!         self.message = Some(message.to_owned());
49//!         Ok(())
50//!     }
51//!
52//!     fn build(self) -> Result<Self::Snapshot, Self::Error> {
53//!         let message = self.message.ok_or_else(|| "No message was assigned.".to_owned())?;
54//!         Ok(Example { message })
55//!     }
56//! }
57//!
58//! /// Subscribes once where the application is wired together, and hands the subscription to its consumer.
59//! fn wire(rc_client: &RemoteConfigurationClient) -> datadog_agent_remote_config::Result<()> {
60//!     let subscription = rc_client.subscribe::<ExampleDecoder>()?;
61//!     tokio::spawn(consume(subscription));
62//!     Ok(())
63//! }
64//!
65//! async fn consume(mut subscription: Subscription<Example>) {
66//!     // A snapshot accepted before the subscription was created is not announced by `changed`, so read it first.
67//!     if let Some(example) = subscription.current() {
68//!         println!("{}", example.message);
69//!     }
70//!     loop {
71//!         match subscription.changed().await {
72//!             Ok(example) => println!("{}", example.message),
73//!             // The client reports the rejection to the Agent; `current` still returns the last accepted snapshot.
74//!             Err(error) => eprintln!("{error}"),
75//!         }
76//!     }
77//! }
78//!
79//! # #[tokio::main(flavor = "current_thread")]
80//! # async fn main() {
81//! // In a test, `TestPublisher` runs the decoder over payloads exactly as the client does.
82//! let (publisher, mut subscription) = TestPublisher::<Example>::new();
83//!
84//! publisher.assign::<ExampleDecoder>([("greeting.v1", "hello"), ("greeting.v2", "hi")]);
85//! assert_eq!(subscription.changed().await.unwrap().message, "hi");
86//!
87//! // A rejection keeps the last accepted snapshot. This decoder rejects an empty assignment only to show that; a real
88//! // decoder must build from one, or an expired Agent cache leaves its last snapshot live.
89//! publisher.assign::<ExampleDecoder>(Vec::<(&str, &str)>::new());
90//! assert_eq!(*subscription.changed().await.unwrap_err(), "No message was assigned.");
91//! assert_eq!(subscription.current().unwrap().message, "hi");
92//! # }
93//! ```
94
95#![deny(missing_docs)]
96
97use std::fmt;
98use std::sync::Arc;
99use std::time::Duration;
100
101use datadog_agent_commons::ipc::client::RemoteAgentClient;
102
103mod decoder;
104mod error;
105mod identity;
106mod json;
107mod metrics;
108mod product;
109mod protocol;
110mod registry;
111mod repository;
112mod source;
113mod subscription;
114#[cfg(any(test, feature = "test-util"))]
115mod test_util;
116#[cfg(test)]
117mod tests;
118mod worker;
119
120pub use decoder::ProductDecoder;
121pub use error::{ApplyError, Error, Result};
122pub use identity::{AgentIdentity, ClientKind};
123pub use json::{decode_json, JsonError};
124pub use product::ConfigId;
125pub use subscription::Subscription;
126#[cfg(any(test, feature = "test-util"))]
127pub use test_util::TestPublisher;
128pub use worker::RemoteConfigurationWorker;
129
130const POLL_INTERVAL: Duration = Duration::from_secs(5);
131const MAX_BACKOFF: Duration = Duration::from_secs(90);
132const REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
133
134/// Settings for a [`RemoteConfigurationClient`].
135///
136/// `Rc` abbreviates Remote Configuration, which avoids the repetition in `RemoteConfigurationClientConfiguration`.
137///
138/// Start from [`new`](Self::new), which takes the kind of client to report and uses the standard schedule, and change
139/// fields as needed. [`validate`](Self::validate) checks the values, and [`RemoteConfigurationClient::new`]
140/// rejects invalid settings.
141///
142/// Whatever the settings, the worker polls once immediately when it starts, and retries every second until its first
143/// successful poll. The exception is an Agent that answers that Remote Configuration is not enabled: the worker then
144/// waits `max_backoff` between attempts, first success or not.
145#[derive(Clone, Debug)]
146pub struct RcClientConfiguration {
147    /// The kind of client the Agent sees, and the details it reports in every poll.
148    ///
149    /// Has no default.
150    pub kind: ClientKind,
151
152    /// How long the worker waits between successful polls.
153    ///
154    /// Shorter intervals deliver changes sooner at the cost of more requests to the Agent. The Agent refreshes from the
155    /// backend far less often than this, so lowering it rarely helps. Must be at least one second.
156    ///
157    /// Defaults to 5 seconds.
158    pub poll_interval: Duration,
159
160    /// The longest the worker waits between polls while polls are failing.
161    ///
162    /// After consecutive failures the wait doubles, with jitter, from `poll_interval` up to this ceiling, and resets
163    /// on the next success. The worker also waits this long between attempts when the Agent has remote configuration
164    /// disabled. Must be at least `poll_interval`.
165    ///
166    /// Defaults to 90 seconds.
167    pub max_backoff: Duration,
168
169    /// How long the worker waits for the Agent to answer one poll.
170    ///
171    /// A poll that runs out of time is abandoned and handled like any other failed poll, so a stuck connection cannot
172    /// stall delivery indefinitely. The Agent may hold the first poll from a new client for up to about two seconds
173    /// while it fetches from the backend, so this must stay well above that. Must be at least one second.
174    ///
175    /// Defaults to 30 seconds.
176    pub request_timeout: Duration,
177}
178
179impl RcClientConfiguration {
180    /// Creates settings for a client that reports itself as `kind`, with the standard schedule.
181    pub fn new(kind: ClientKind) -> Self {
182        Self {
183            kind,
184            poll_interval: POLL_INTERVAL,
185            max_backoff: MAX_BACKOFF,
186            request_timeout: REQUEST_TIMEOUT,
187        }
188    }
189
190    /// Checks that the settings are valid.
191    ///
192    /// # Errors
193    ///
194    /// Returns [`Error::InvalidSettings`] if an agent's [`name`](AgentIdentity::name) or
195    /// [`version`](AgentIdentity::version) is empty, `poll_interval` is shorter than one second, `max_backoff` is
196    /// shorter than `poll_interval`, or `request_timeout` is shorter than one second. The fields are checked in that
197    /// order, and the error describes the first that fails.
198    pub fn validate(&self) -> Result<()> {
199        let ClientKind::Agent(agent) = &self.kind;
200        if agent.name.is_empty() {
201            return Err(Error::InvalidSettings {
202                message: "The agent name must not be empty.".to_owned(),
203            });
204        }
205        if agent.version.is_empty() {
206            return Err(Error::InvalidSettings {
207                message: "The agent version must not be empty.".to_owned(),
208            });
209        }
210        if self.poll_interval < Duration::from_secs(1) {
211            return Err(Error::InvalidSettings {
212                message: format!(
213                    "poll_interval must be at least one second; got {:?}.",
214                    self.poll_interval
215                ),
216            });
217        }
218        if self.max_backoff < self.poll_interval {
219            return Err(Error::InvalidSettings {
220                message: format!(
221                    "max_backoff must be at least poll_interval ({:?}); got {:?}.",
222                    self.poll_interval, self.max_backoff
223                ),
224            });
225        }
226        if self.request_timeout < Duration::from_secs(1) {
227            return Err(Error::InvalidSettings {
228                message: format!(
229                    "request_timeout must be at least one second; got {:?}.",
230                    self.request_timeout
231                ),
232            });
233        }
234        // Settings are valid.
235        Ok(())
236    }
237}
238
239/// A cloneable handle for subscribing to Remote Configuration products.
240///
241/// Every clone shares one set of subscriptions and one client identity with the worker.
242#[derive(Clone)]
243pub struct RemoteConfigurationClient {
244    shared: Arc<registry::Shared>,
245}
246
247impl fmt::Debug for RemoteConfigurationClient {
248    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
249        f.debug_struct("RemoteConfigurationClient")
250            .field("client_id", &self.shared.client_id)
251            .finish_non_exhaustive()
252    }
253}
254
255// TODO: consider opt-in health notifications when a subscriber needs them (e.g. CWS enforcement).
256// TODO: consider per-product option to keep last good when the Agent reports expired (e.g. Cluster Agent autoscaling).
257impl RemoteConfigurationClient {
258    /// Creates a client and its worker from a connected Datadog Agent client.
259    ///
260    /// The connection must be dedicated to Remote Configuration rather than a clone of another subsystem's client,
261    /// following the Saluki pattern of one connection per subsystem. The connection type is cloneable, so nothing here
262    /// enforces that; it is a convention of the caller's wiring. Construction does not spawn the worker or probe
263    /// Remote Configuration availability; the caller schedules the worker through a supervisor or its `run` method.
264    ///
265    /// The client identifies itself to the Agent with a random ID, generated here and kept for the life of the client,
266    /// and with the kind and details in `config`.
267    ///
268    /// # Errors
269    ///
270    /// Returns an error if `config` fails [`RcClientConfiguration::validate`].
271    pub fn new(ra: RemoteAgentClient, config: RcClientConfiguration) -> Result<(Self, RemoteConfigurationWorker)> {
272        config.validate()?;
273        Ok(Self::with_agent(Box::new(ra), config))
274    }
275
276    /// Creates a client and its worker around any [`RcAgent`](source::RcAgent), so tests can replace the Agent.
277    pub(crate) fn with_agent(
278        agent: Box<dyn source::RcAgent>, config: RcClientConfiguration,
279    ) -> (Self, RemoteConfigurationWorker) {
280        let shared = Arc::new(registry::Shared::new());
281        let worker = RemoteConfigurationWorker::new(Arc::clone(&shared), agent, config);
282        (Self { shared }, worker)
283    }
284
285    /// Subscribes to the product `P` decodes, [`P::PRODUCT`](ProductDecoder::PRODUCT).
286    ///
287    /// The decoder is named here, where how a product is read is the subject; the returned subscription is typed by the
288    /// snapshot that decoder builds. Subscriptions can be added while the worker runs.
289    ///
290    /// A product may have only one live subscription per client, even when two decoders name it. Several consumers of
291    /// one product therefore share a single [`Subscription`] by cloning it, rather than each subscribing for
292    /// themselves. Dropping the last clone unsubscribes, after which the product may be subscribed again.
293    /// The new subscription starts with no accepted snapshot.
294    ///
295    /// Subscribing wakes the worker to poll without waiting for its next scheduled poll. The Agent answers polls from a
296    /// cache that it refreshes from the Datadog backend on its own schedule, every minute by default. It refreshes
297    /// early only for a client it has not seen before, so a product that no other client of the Agent requests,
298    /// subscribed after this client's first poll, may take up to the Agent's refresh interval to arrive.
299    ///
300    /// # Errors
301    ///
302    /// Returns [`Error::AlreadySubscribed`] when the product still has a live subscription on this client, which
303    /// indicates that the caller should be receiving a clone of the existing subscription instead.
304    pub fn subscribe<P>(&self) -> Result<Subscription<P::Snapshot, P::Error>>
305    where
306        P: ProductDecoder,
307    {
308        self.shared.subscribe::<P>()
309    }
310}