aster_forge_config/notification/
message.rs1use serde::{Deserialize, Serialize};
2use std::future::Future;
3
4use crate::Result;
5
6#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
8#[serde(rename_all = "snake_case")]
9pub enum ConfigNotificationSource {
10 Api,
12 Cli,
14 Startup,
16 Other(String),
18}
19
20#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
22pub struct ConfigReloadMessage {
23 pub namespace: String,
25 pub origin_runtime_id: String,
28 pub keys: Vec<String>,
30 pub source: ConfigNotificationSource,
32}
33
34impl ConfigReloadMessage {
35 pub fn new(
37 namespace: impl Into<String>,
38 origin_runtime_id: impl Into<String>,
39 keys: impl IntoIterator<Item = impl Into<String>>,
40 source: ConfigNotificationSource,
41 ) -> Self {
42 let mut keys = keys.into_iter().map(Into::into).collect::<Vec<_>>();
43 keys.sort();
44 keys.dedup();
45 Self {
46 namespace: namespace.into(),
47 origin_runtime_id: origin_runtime_id.into(),
48 keys,
49 source,
50 }
51 }
52
53 pub fn encode(&self) -> Result<String> {
59 serde_json::to_string(self).map_err(Into::into)
60 }
61
62 pub fn decode(payload: &str) -> Result<Self> {
68 serde_json::from_str(payload).map_err(Into::into)
69 }
70}
71
72#[derive(Debug, Clone, PartialEq, Eq)]
74pub enum ConfigChangeEvent {
75 Reload(ConfigReloadMessage),
77}
78
79impl ConfigChangeEvent {
80 #[must_use]
82 pub const fn reload_message(&self) -> &ConfigReloadMessage {
83 match self {
84 Self::Reload(message) => message,
85 }
86 }
87}
88
89#[derive(Debug, Clone, Copy, PartialEq, Eq)]
91pub enum ConfigReloadDecision {
92 Reloaded,
94 IgnoredNamespace,
96 IgnoredOrigin,
98}
99
100impl ConfigReloadDecision {
101 #[must_use]
103 pub const fn as_label(self) -> &'static str {
104 match self {
105 Self::Reloaded => "reloaded",
106 Self::IgnoredNamespace => "ignored_namespace",
107 Self::IgnoredOrigin => "ignored_origin",
108 }
109 }
110}
111
112#[derive(Debug, Clone, PartialEq, Eq)]
114pub struct ConfigReloadWorkerConfig {
115 pub namespace: String,
117 pub runtime_id: String,
119}
120
121impl ConfigReloadWorkerConfig {
122 pub fn new(namespace: impl Into<String>, runtime_id: impl Into<String>) -> Self {
124 Self {
125 namespace: namespace.into(),
126 runtime_id: runtime_id.into(),
127 }
128 }
129
130 #[must_use]
132 pub fn accepts_namespace(&self, message: &ConfigReloadMessage) -> bool {
133 message.namespace == self.namespace
134 }
135
136 #[must_use]
138 pub fn is_local_origin(&self, message: &ConfigReloadMessage) -> bool {
139 message.origin_runtime_id == self.runtime_id
140 }
141}
142
143pub fn decode_config_reload_transport_payload(payload: &str) -> Result<ConfigChangeEvent> {
153 ConfigReloadMessage::decode(payload).map(ConfigChangeEvent::Reload)
154}
155
156pub async fn handle_config_reload_notification<F, Fut>(
162 config: &ConfigReloadWorkerConfig,
163 message: ConfigReloadMessage,
164 reload: F,
165) -> Result<ConfigReloadDecision>
166where
167 F: FnOnce(ConfigReloadMessage) -> Fut,
168 Fut: Future<Output = Result<()>>,
169{
170 if !config.accepts_namespace(&message) {
171 return Ok(ConfigReloadDecision::IgnoredNamespace);
172 }
173 if config.is_local_origin(&message) {
174 return Ok(ConfigReloadDecision::IgnoredOrigin);
175 }
176
177 reload(message).await?;
178 Ok(ConfigReloadDecision::Reloaded)
179}