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> {
55 serde_json::to_string(self).map_err(Into::into)
56 }
57
58 pub fn decode(payload: &str) -> Result<Self> {
60 serde_json::from_str(payload).map_err(Into::into)
61 }
62}
63
64#[derive(Debug, Clone, PartialEq, Eq)]
66pub enum ConfigChangeEvent {
67 Reload(ConfigReloadMessage),
69}
70
71impl ConfigChangeEvent {
72 pub const fn reload_message(&self) -> &ConfigReloadMessage {
74 match self {
75 Self::Reload(message) => message,
76 }
77 }
78}
79
80#[derive(Debug, Clone, Copy, PartialEq, Eq)]
82pub enum ConfigReloadDecision {
83 Reloaded,
85 IgnoredNamespace,
87 IgnoredOrigin,
89}
90
91impl ConfigReloadDecision {
92 pub const fn as_label(self) -> &'static str {
94 match self {
95 Self::Reloaded => "reloaded",
96 Self::IgnoredNamespace => "ignored_namespace",
97 Self::IgnoredOrigin => "ignored_origin",
98 }
99 }
100}
101
102#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct ConfigReloadWorkerConfig {
105 pub namespace: String,
107 pub runtime_id: String,
109}
110
111impl ConfigReloadWorkerConfig {
112 pub fn new(namespace: impl Into<String>, runtime_id: impl Into<String>) -> Self {
114 Self {
115 namespace: namespace.into(),
116 runtime_id: runtime_id.into(),
117 }
118 }
119
120 pub fn accepts_namespace(&self, message: &ConfigReloadMessage) -> bool {
122 message.namespace == self.namespace
123 }
124
125 pub fn is_local_origin(&self, message: &ConfigReloadMessage) -> bool {
127 message.origin_runtime_id == self.runtime_id
128 }
129}
130
131pub fn decode_config_reload_transport_payload(payload: &str) -> Result<ConfigChangeEvent> {
137 ConfigReloadMessage::decode(payload).map(ConfigChangeEvent::Reload)
138}
139
140pub async fn handle_config_reload_notification<F, Fut>(
142 config: &ConfigReloadWorkerConfig,
143 message: ConfigReloadMessage,
144 reload: F,
145) -> Result<ConfigReloadDecision>
146where
147 F: FnOnce(ConfigReloadMessage) -> Fut,
148 Fut: Future<Output = Result<()>>,
149{
150 if !config.accepts_namespace(&message) {
151 return Ok(ConfigReloadDecision::IgnoredNamespace);
152 }
153 if config.is_local_origin(&message) {
154 return Ok(ConfigReloadDecision::IgnoredOrigin);
155 }
156
157 reload(message).await?;
158 Ok(ConfigReloadDecision::Reloaded)
159}