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}