datadog_agent_remote_config/
subscription.rs

1//! Typed subscriptions to a product's configuration.
2
3use std::fmt;
4use std::sync::Arc;
5
6use tokio::sync::watch;
7
8use crate::decoder::Outcome;
9
10/// The published state of one product.
11///
12/// A rejection never evicts `accepted`, so retaining the last known-good configuration is the client's behavior rather
13/// than each consumer's bookkeeping.
14pub(crate) struct Snapshot<T, E> {
15    /// The most recently accepted configuration, absent until the first snapshot is accepted.
16    pub(crate) accepted: Option<Arc<T>>,
17
18    /// The rejection of the most recent snapshot, absent while the most recent snapshot was accepted.
19    pub(crate) rejection: Option<Arc<E>>,
20}
21
22/// A handle that observes one product's configuration as the client publishes it.
23///
24/// `T` is the snapshot the product's [`ProductDecoder`](crate::ProductDecoder) builds and `E` is the error it produces,
25/// so a subscription reads as the value it delivers rather than as the decoder that produced it. A product with nothing
26/// richer to report leaves `E` at its default of [`String`].
27///
28/// A consumer reads [`current`](Self::current) and then loops on [`changed`](Self::changed). A subscription created
29/// after the client has already published treats that value as observed and receives no notification for it, so a
30/// consumer that only awaited `changed` would wait for a snapshot that may never arrive.
31///
32/// Cloning shares one subscription between several consumers: each clone tracks its own position, while decoding
33/// happens once per snapshot. Slow consumers may skip intermediate publications and observe only the latest state.
34/// Dropping the last clone unsubscribes the product.
35pub struct Subscription<T, E = String> {
36    pub(crate) receiver: watch::Receiver<Snapshot<T, E>>,
37}
38
39impl<T, E> Subscription<T, E> {
40    /// Creates a subscription that never receives a snapshot.
41    ///
42    /// Use this where a component takes a subscription but the process has no client to subscribe with, for example
43    /// because Remote Configuration is turned off. The component then runs the same code path as one whose client has
44    /// not delivered anything yet: [`current`](Self::current) always returns `None`, and [`changed`](Self::changed)
45    /// waits indefinitely.
46    ///
47    /// An inert subscription is not registered with any client, so it does not count against a client's one
48    /// subscription per product.
49    pub fn inert() -> Self {
50        // `changed` turns a closed channel into a future that never completes.
51        let (publisher, subscription) = Publisher::new();
52        drop(publisher);
53        subscription
54    }
55
56    /// Returns the most recently accepted configuration.
57    ///
58    /// Returns `None` until the first snapshot is accepted, whether because the worker has not yet polled, is not
59    /// running, or the subscription is [inert](Self::inert). `None` does not mean the product has no configuration
60    /// assigned: a fresh decoder's [`build`](crate::ProductDecoder::build) must produce a snapshot even for an empty
61    /// assignment.
62    pub fn current(&self) -> Option<Arc<T>> {
63        self.receiver.borrow().accepted.clone()
64    }
65
66    /// Waits for a newly published snapshot.
67    ///
68    /// This takes `&mut self` because the subscription tracks which publication it has observed. Each task that waits
69    /// holds its own clone.
70    ///
71    /// Slow consumers may skip intermediate publications and observe only the latest state. Once the worker stops,
72    /// this waits indefinitely after any pending publication has been observed, so a caller may `select!` on it.
73    ///
74    /// A published snapshot is not guaranteed to differ from the previous one: after the worker restarts, it decodes
75    /// every product again.
76    ///
77    /// # Errors
78    ///
79    /// Returns the subscriber's own decoding error when the snapshot was rejected. The client reports that rejection to
80    /// the Agent independently; a consumer with nothing to report can ignore it. [`current`](Self::current) continues
81    /// to return the last accepted configuration.
82    pub async fn changed(&mut self) -> Result<Arc<T>, Arc<E>> {
83        if self.receiver.changed().await.is_err() {
84            return std::future::pending().await;
85        }
86
87        let snapshot = self.receiver.borrow_and_update();
88        if let Some(error) = &snapshot.rejection {
89            Err(Arc::clone(error))
90        } else {
91            Ok(Arc::clone(
92                snapshot.accepted.as_ref().expect("published snapshot has a value"),
93            ))
94        }
95    }
96}
97
98impl<T, E> fmt::Debug for Subscription<T, E> {
99    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
100        let accepted = self.receiver.borrow().accepted.is_some();
101        f.debug_struct("Subscription")
102            .field("accepted", &accepted)
103            .finish_non_exhaustive()
104    }
105}
106
107impl<T, E> Clone for Subscription<T, E> {
108    fn clone(&self) -> Self {
109        Self {
110            receiver: self.receiver.clone(),
111        }
112    }
113}
114
115/// The publishing side of one product's [`Subscription`] and its clones.
116///
117/// The client's registry and [`TestPublisher`](crate::TestPublisher) both publish through this, so a subscription
118/// behaves the same in a subscriber's test as in production.
119pub(crate) struct Publisher<T, E> {
120    sender: watch::Sender<Snapshot<T, E>>,
121}
122
123impl<T, E> Publisher<T, E> {
124    /// Creates a publisher and a subscription with no accepted snapshot.
125    pub(crate) fn new() -> (Self, Subscription<T, E>) {
126        let (sender, receiver) = watch::channel(Snapshot {
127            accepted: None,
128            rejection: None,
129        });
130        (Self { sender }, Subscription { receiver })
131    }
132
133    pub(crate) fn accept(&self, snapshot: T) {
134        // The replaced values are dropped after `send_modify` releases the channel's lock, so a subscriber's `Drop` for
135        // its previous snapshot or rejection can read the subscription instead of deadlocking. A panic in that `Drop`
136        // escapes the decoder's `catch_unwind` and ends the worker, which its supervisor restarts.
137        let mut replaced = (None, None);
138        self.sender.send_modify(|state| {
139            replaced = (state.accepted.replace(Arc::new(snapshot)), state.rejection.take());
140        });
141        drop(replaced);
142    }
143
144    pub(crate) fn reject(&self, error: E) {
145        // Dropped after the lock is released, as in `accept`.
146        let mut replaced = None;
147        self.sender.send_modify(|state| {
148            replaced = state.rejection.replace(Arc::new(error));
149        });
150        drop(replaced);
151    }
152
153    /// Publishes the outcome of a decoder run. A panicked run publishes nothing, so subscribers are not notified.
154    pub(crate) fn publish(&self, outcome: Outcome<T, E>) {
155        match outcome {
156            Outcome::Accepted(snapshot) => self.accept(snapshot),
157            Outcome::Rejected { error, .. } => self.reject(error),
158            Outcome::Panicked => {}
159        }
160    }
161
162    /// Returns whether any clone of the subscription is still alive.
163    ///
164    /// Once this is false it stays false: a subscription can only be cloned from a live one.
165    pub(crate) fn is_subscribed(&self) -> bool {
166        self.sender.receiver_count() > 0
167    }
168}