datadog_agent_remote_config/
decoder.rs

1//! The subscriber's decoding contract.
2
3use std::panic::{catch_unwind, AssertUnwindSafe};
4
5use crate::{ApplyError, ConfigId};
6
7/// Decodes one product's assigned configurations into a snapshot.
8///
9/// A decoder names the product it decodes, so subscribing with a decoder cannot attach it to the wrong product. Define
10/// a decoder in the crate that owns its snapshot type. A decoder for a product whose type has no other owner belongs in
11/// this crate instead; otherwise this crate knows no product's payload.
12///
13/// The client drives an implementation once per published snapshot: it constructs a decoder with [`Default`], calls
14/// [`decode`](Self::decode) once per assigned configuration in ascending [`ConfigId`] order, and then calls
15/// [`build`](Self::build). Ascending order is part of the contract, so a reduction that keeps the last valid
16/// configuration is well defined.
17///
18/// A product is decoded again only when its assigned configurations or their contents change, so a rejected
19/// assignment is not retried until it changes. When the Agent reports its configuration expired, every product that
20/// had configurations is decoded as an empty assignment; a product that already had none is not decoded again. The
21/// client then asks the Agent for the full assignment until a response is not expired, so configurations that return
22/// are decoded as new.
23///
24/// A rejected snapshot leaves subscribers with the last accepted one, so an implementation must build a snapshot
25/// from an empty assignment; see [`build`](Self::build).
26///
27/// Decoding runs on the client's worker task and delays polling for every product while it runs, so implementations
28/// **MUST NOT** block. A panic in [`decode`](Self::decode) or [`build`](Self::build) is caught: the client discards the
29/// decoder, rejects every configuration in the assignment, and publishes nothing, so subscribers keep the last accepted
30/// snapshot and are not notified.
31///
32/// # Design
33///
34/// A product's assignment is several configurations, each with its own payload, and the protocol carries an apply
35/// status for each one separately. Decoding therefore accumulates configuration by configuration rather than consuming
36/// the assignment as a whole: a configuration that [`decode`](Self::decode) rejects is attributed to itself and
37/// skipped, and the decoder still builds from the rest. Because the client drives the loop, that attribution cannot be
38/// forgotten or misdirected. Whether a rejected configuration invalidates the whole snapshot remains the decoder's
39/// choice: [`build`](Self::build) may reject the snapshot, for example when a configuration it requires was rejected.
40///
41/// Accumulating into the decoder, rather than returning one decoded value per configuration, is what lets a product
42/// whose configurations have different shapes hold each of them in a field of its own type instead of funneling them
43/// through a shared enum.
44///
45/// # Examples
46///
47/// A product that accepts at most one configuration records every configuration it is given and rejects more than one
48/// in [`build`](Self::build), since only `build` sees the whole assignment:
49///
50/// ```
51/// use datadog_agent_remote_config::{ConfigId, ProductDecoder};
52///
53/// /// Holds the product's single configuration, if it has one.
54/// #[derive(Default)]
55/// struct SingleDecoder {
56///     payloads: Vec<(ConfigId, Vec<u8>)>,
57/// }
58///
59/// impl ProductDecoder for SingleDecoder {
60///     const PRODUCT: &'static str = "EXAMPLE_SINGLE";
61///
62///     type Snapshot = Option<Vec<u8>>;
63///     type Error = String;
64///
65///     fn decode(&mut self, id: &ConfigId, payload: &[u8]) -> Result<(), Self::Error> {
66///         self.payloads.push((id.clone(), payload.to_vec()));
67///         Ok(())
68///     }
69///
70///     fn build(mut self) -> Result<Self::Snapshot, Self::Error> {
71///         if self.payloads.len() > 1 {
72///             return Err(format!("Expected at most one configuration; got {}.", self.payloads.len()));
73///         }
74///         Ok(self.payloads.pop().map(|(_, payload)| payload))
75///     }
76/// }
77///
78/// let mut decoder = SingleDecoder::default();
79/// decoder.decode(&ConfigId::new("a"), b"one").unwrap();
80/// decoder.decode(&ConfigId::new("b"), b"two").unwrap();
81/// assert!(decoder.build().is_err());
82/// ```
83pub trait ProductDecoder: Default + Send + 'static {
84    /// The product this decoder decodes, by its protocol name, such as `APM_SEMANTIC_CORE_DD`.
85    ///
86    /// The client requests this product from the Agent and keys the product's subscription by this name. The client
87    /// can only subscribe once per product: attempting a second subscription to the same product while the first is
88    /// live returns [`Error::AlreadySubscribed`](crate::Error::AlreadySubscribed), whichever decoder it uses.
89    const PRODUCT: &'static str;
90
91    /// The type that holds the product's decoded configuration.
92    ///
93    /// [`build`](Self::build) produces it from the product's configurations, and subscribers receive it.
94    type Snapshot: Send + Sync + 'static;
95
96    /// The error this product's decoding and validation produces.
97    ///
98    /// Use [`String`] when there is nothing richer to report; a product that wants to describe a failure more
99    /// precisely for its own diagnostics defines its own type and implements [`ApplyError`] for it.
100    type Error: ApplyError + Send + Sync + 'static;
101
102    /// Accumulates one of the product's assigned configurations.
103    ///
104    /// The client calls this at most once per distinct `id` within a single snapshot, then calls
105    /// [`build`](Self::build) once. A decoder may therefore hold one slot per configuration ID without a second
106    /// configuration silently displacing the first.
107    ///
108    /// # Errors
109    ///
110    /// Returns [`Self::Error`] to reject this configuration alone. The client attributes the rejection to it and skips
111    /// it, and still decodes the product's remaining configurations.
112    fn decode(&mut self, id: &ConfigId, payload: &[u8]) -> Result<(), Self::Error>;
113
114    /// Validates the accumulated configurations and produces the snapshot to publish.
115    ///
116    /// Check requirements across configurations here, including whether a required configuration is present when the
117    /// decoder receives a non-empty assignment. The client also calls this without calling [`decode`](Self::decode)
118    /// when no configurations are assigned (after the backend removes the last one or the Agent reports its
119    /// configuration expired), or when every assigned configuration was rejected for sharing an ID.
120    ///
121    /// An empty assignment to the decoder means "no configuration," even if files were assigned but all collided.
122    /// This method **MUST** build a snapshot representing that state. Rejecting an empty assignment leaves subscribers
123    /// reading configuration the backend removed or that expired with the Agent's cache until a non-empty assignment
124    /// builds successfully, because a rejection keeps the last accepted snapshot.
125    ///
126    /// The client acknowledges successfully decoded configurations only when this method succeeds.
127    ///
128    /// # Errors
129    ///
130    /// Returns [`Self::Error`] to reject the snapshot. The client rejects every successfully decoded configuration
131    /// with this error's apply reason; configurations rejected by [`decode`](Self::decode) keep their own errors.
132    /// No new snapshot is published, and subscribers retain the last accepted snapshot.
133    ///
134    /// This method cannot reject selected configurations while publishing the rest. Selective rejection belongs in
135    /// [`decode`](Self::decode).
136    fn build(self) -> Result<Self::Snapshot, Self::Error>;
137}
138
139/// The fixed reason reported for every configuration when a decoder panics.
140///
141/// A panic is a bug in the decoder rather than a verdict about any one configuration's bytes, so it overwrites any
142/// rejection a configuration had already earned, including from an earlier `decode` call.
143pub(crate) const PANICKED: &str = "Product decoder panicked.";
144
145/// What one run of a decoder over a product's assignment produced.
146pub(crate) struct Evaluation<T, E> {
147    pub(crate) outcome: Outcome<T, E>,
148    /// Each assigned configuration's verdict, in ascending ID order.
149    pub(crate) verdicts: Vec<(ConfigId, Verdict)>,
150}
151
152/// The verdict one configuration earned from a run of a decoder.
153///
154/// The reason is the string the client reports to the Agent; the variant is which rule produced it, which is what
155/// the client counts.
156#[derive(Debug, PartialEq, Eq)]
157pub(crate) enum Verdict {
158    /// `decode` accepted this configuration and the snapshot it was part of was built.
159    Acknowledged,
160
161    /// `decode` rejected this configuration, with the reason it gave.
162    DecodeRejected(String),
163
164    /// `decode` accepted this configuration, but `build` rejected the snapshot, with the reason it gave.
165    BuildRejected(String),
166
167    /// The decoder panicked, so this configuration is rejected without a reason of its own.
168    Panicked,
169
170    /// This configuration shares its ID with another, so decode did not run.
171    ///
172    /// Set by the client rather than a decoder; see [`Repository`](crate::repository::Repository).
173    Collided(String),
174}
175
176pub(crate) enum Outcome<T, E> {
177    /// `build` succeeded; the snapshot is published.
178    Accepted(T),
179
180    /// `build` failed; the error is published as a rejection and `current` keeps the last accepted snapshot.
181    Rejected {
182        error: E,
183
184        /// The error's apply reason, computed inside the decoder's `catch_unwind` so that a panic in
185        /// [`ApplyError::apply_error`] is caught like any other decoder panic.
186        reason: String,
187    },
188
189    /// `decode` or `build` panicked; nothing is published and subscribers are not notified.
190    Panicked,
191}
192
193/// Runs a fresh decoder over one product's assignment exactly as the client does.
194///
195/// This is the only implementation of the decoding rules, shared by the worker and by
196/// [`TestPublisher::assign`](crate::TestPublisher::assign), so that what a subscriber tests is what production runs:
197/// configurations in ascending [`ConfigId`] order, a rejected configuration skipped while the rest are still decoded,
198/// then `build`, with a panic in either caught.
199pub(crate) fn evaluate<P: ProductDecoder>(mut assignment: Vec<(ConfigId, &[u8])>) -> Evaluation<P::Snapshot, P::Error> {
200    assignment.sort_by(|(left, _), (right, _)| left.cmp(right));
201
202    let mut verdicts = Vec::with_capacity(assignment.len());
203    let result = catch_unwind(AssertUnwindSafe(|| {
204        let mut decoder = P::default();
205        for (id, payload) in &assignment {
206            let verdict = match decoder.decode(id, payload) {
207                Ok(()) => Verdict::Acknowledged,
208                Err(error) => Verdict::DecodeRejected(error.apply_error()),
209            };
210            verdicts.push((id.clone(), verdict));
211        }
212        match decoder.build() {
213            Ok(snapshot) => Outcome::Accepted(snapshot),
214            Err(error) => {
215                let reason = error.apply_error();
216                for (_, verdict) in &mut verdicts {
217                    if matches!(verdict, Verdict::Acknowledged) {
218                        *verdict = Verdict::BuildRejected(reason.clone());
219                    }
220                }
221                Outcome::Rejected { error, reason }
222            }
223        }
224    }));
225
226    let outcome = match result {
227        Ok(outcome) => outcome,
228        Err(_) => {
229            verdicts = assignment.into_iter().map(|(id, _)| (id, Verdict::Panicked)).collect();
230            Outcome::Panicked
231        }
232    };
233
234    Evaluation { outcome, verdicts }
235}