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}