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}