1use super::loader::{
16 MergedTargetConfigRecord, collect_merged_target_configs_compat_from_env, collect_merged_target_configs_from_env,
17};
18use crate::TargetError;
19use crate::domain::TargetDomain;
20use rustfs_config::server_config::{Config, KVS};
21use std::collections::HashSet;
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24pub struct TargetPluginInstanceCompatDescriptor<'a> {
25 pub domain: TargetDomain,
26 pub plugin_id: &'a str,
27 pub target_type: &'a str,
28 pub subsystem: &'a str,
29 pub route_prefix: &'a str,
30 pub valid_fields: &'a [&'a str],
31}
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq)]
34pub enum TargetInstanceSourceClass {
35 Config,
36 Env,
37 Mixed,
38}
39
40#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
41pub struct TargetInstanceSourceHints {
42 pub has_file_default: bool,
43 pub has_file_instance: bool,
44 pub has_env_default: bool,
45 pub has_env_instance: bool,
46}
47
48impl TargetInstanceSourceHints {
49 #[inline]
50 pub fn has_config_source(self) -> bool {
51 self.has_file_default || self.has_file_instance
52 }
53
54 #[inline]
55 pub fn has_env_source(self) -> bool {
56 self.has_env_default || self.has_env_instance
57 }
58
59 #[inline]
60 pub fn classification(self) -> TargetInstanceSourceClass {
61 match (self.has_config_source(), self.has_env_source()) {
62 (true, true) => TargetInstanceSourceClass::Mixed,
63 (true, false) => TargetInstanceSourceClass::Config,
64 (false, true) => TargetInstanceSourceClass::Env,
65 (false, false) => TargetInstanceSourceClass::Config,
66 }
67 }
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct TargetPluginInstanceRecord {
72 pub domain: TargetDomain,
73 pub plugin_id: String,
74 pub target_type: String,
75 pub subsystem: String,
76 pub instance_id: String,
77 pub enabled: bool,
78 pub source_hints: TargetInstanceSourceHints,
79 pub effective_config: KVS,
80}
81
82pub type LegacyTargetInstanceDescriptor<'a> = TargetPluginInstanceCompatDescriptor<'a>;
83pub type TargetPluginInstance = TargetPluginInstanceRecord;
84
85pub fn normalize_target_plugin_instances(
86 config: &Config,
87 descriptor: &TargetPluginInstanceCompatDescriptor<'_>,
88) -> Vec<TargetPluginInstanceRecord> {
89 normalize_target_plugin_instances_from_env(config, descriptor, std::env::vars())
90}
91
92pub fn try_normalize_target_plugin_instances(
93 config: &Config,
94 descriptor: &TargetPluginInstanceCompatDescriptor<'_>,
95) -> Result<Vec<TargetPluginInstanceRecord>, TargetError> {
96 try_normalize_target_plugin_instances_from_env(config, descriptor, std::env::vars())
97}
98
99pub fn normalize_target_plugin_instances_from_env<I>(
100 config: &Config,
101 descriptor: &TargetPluginInstanceCompatDescriptor<'_>,
102 env_vars: I,
103) -> Vec<TargetPluginInstanceRecord>
104where
105 I: IntoIterator<Item = (String, String)>,
106{
107 let valid_fields = descriptor
108 .valid_fields
109 .iter()
110 .map(|field| (*field).to_string())
111 .collect::<HashSet<_>>();
112
113 collect_merged_target_configs_compat_from_env(
114 config,
115 descriptor.subsystem,
116 descriptor.route_prefix,
117 descriptor.target_type,
118 &valid_fields,
119 env_vars,
120 )
121 .into_iter()
122 .map(|record| target_plugin_instance_record(descriptor, record))
123 .collect()
124}
125
126pub fn try_normalize_target_plugin_instances_from_env<I>(
127 config: &Config,
128 descriptor: &TargetPluginInstanceCompatDescriptor<'_>,
129 env_vars: I,
130) -> Result<Vec<TargetPluginInstanceRecord>, TargetError>
131where
132 I: IntoIterator<Item = (String, String)>,
133{
134 let valid_fields = descriptor
135 .valid_fields
136 .iter()
137 .map(|field| (*field).to_string())
138 .collect::<HashSet<_>>();
139
140 Ok(collect_merged_target_configs_from_env(
141 config,
142 descriptor.subsystem,
143 descriptor.route_prefix,
144 descriptor.target_type,
145 &valid_fields,
146 env_vars,
147 )?
148 .into_iter()
149 .map(|record| target_plugin_instance_record(descriptor, record))
150 .collect())
151}
152
153fn target_plugin_instance_record(
154 descriptor: &TargetPluginInstanceCompatDescriptor<'_>,
155 record: MergedTargetConfigRecord,
156) -> TargetPluginInstanceRecord {
157 TargetPluginInstanceRecord {
158 domain: descriptor.domain,
159 plugin_id: descriptor.plugin_id.to_string(),
160 target_type: descriptor.target_type.to_string(),
161 subsystem: descriptor.subsystem.to_string(),
162 instance_id: record.instance_id,
163 enabled: record.enabled,
164 source_hints: TargetInstanceSourceHints {
165 has_file_default: record.has_file_default,
166 has_file_instance: record.has_file_instance,
167 has_env_default: record.has_env_default,
168 has_env_instance: record.has_env_instance,
169 },
170 effective_config: record.effective_config,
171 }
172}
173
174pub fn normalize_legacy_target_instances(
175 config: &Config,
176 descriptor: &LegacyTargetInstanceDescriptor<'_>,
177) -> Vec<TargetPluginInstance> {
178 normalize_target_plugin_instances(config, descriptor)
179}
180
181pub fn normalize_legacy_target_instances_from_env<I>(
182 config: &Config,
183 descriptor: &LegacyTargetInstanceDescriptor<'_>,
184 env_vars: I,
185) -> Vec<TargetPluginInstance>
186where
187 I: IntoIterator<Item = (String, String)>,
188{
189 normalize_target_plugin_instances_from_env(config, descriptor, env_vars)
190}
191
192#[cfg(test)]
193mod tests {
194 use super::{
195 TargetInstanceSourceClass, TargetPluginInstanceCompatDescriptor, normalize_legacy_target_instances_from_env,
196 try_normalize_target_plugin_instances_from_env,
197 };
198 use crate::TargetError;
199 use crate::domain::TargetDomain;
200 use crate::manifest::builtin_target_manifest;
201 use rustfs_config::audit::{AUDIT_ROUTE_PREFIX, AUDIT_WEBHOOK_KEYS, AUDIT_WEBHOOK_SUB_SYS};
202 use rustfs_config::notify::{NOTIFY_ROUTE_PREFIX, NOTIFY_WEBHOOK_KEYS, NOTIFY_WEBHOOK_SUB_SYS};
203 use rustfs_config::server_config::{Config, KVS};
204 use rustfs_config::{ENABLE_KEY, WEBHOOK_ENDPOINT, WEBHOOK_QUEUE_LIMIT};
205 use std::collections::HashMap;
206
207 fn notify_webhook_descriptor() -> TargetPluginInstanceCompatDescriptor<'static> {
208 TargetPluginInstanceCompatDescriptor {
209 domain: TargetDomain::Notify,
210 plugin_id: builtin_target_manifest("webhook").plugin_id,
211 target_type: "webhook",
212 subsystem: NOTIFY_WEBHOOK_SUB_SYS,
213 route_prefix: NOTIFY_ROUTE_PREFIX,
214 valid_fields: NOTIFY_WEBHOOK_KEYS,
215 }
216 }
217
218 fn audit_webhook_descriptor() -> TargetPluginInstanceCompatDescriptor<'static> {
219 TargetPluginInstanceCompatDescriptor {
220 domain: TargetDomain::Audit,
221 plugin_id: builtin_target_manifest("webhook").plugin_id,
222 target_type: "webhook",
223 subsystem: AUDIT_WEBHOOK_SUB_SYS,
224 route_prefix: AUDIT_ROUTE_PREFIX,
225 valid_fields: AUDIT_WEBHOOK_KEYS,
226 }
227 }
228
229 #[test]
230 fn normalize_notify_instances_merges_file_and_env_sources() {
231 let mut cfg = Config(HashMap::new());
232 let mut subsystem = HashMap::new();
233
234 let mut default_kvs = KVS::new();
235 default_kvs.insert(ENABLE_KEY.to_string(), "on".to_string());
236 default_kvs.insert(WEBHOOK_QUEUE_LIMIT.to_string(), "10".to_string());
237 subsystem.insert("_".to_string(), default_kvs);
238
239 let mut primary = KVS::new();
240 primary.insert(WEBHOOK_ENDPOINT.to_string(), "https://example.com/primary".to_string());
241 subsystem.insert("primary".to_string(), primary);
242
243 cfg.0.insert(NOTIFY_WEBHOOK_SUB_SYS.to_string(), subsystem);
244
245 let instances = normalize_legacy_target_instances_from_env(
246 &cfg,
247 ¬ify_webhook_descriptor(),
248 vec![
249 ("RUSTFS_NOTIFY_WEBHOOK_QUEUE_LIMIT".to_string(), "42".to_string()),
250 ("RUSTFS_NOTIFY_WEBHOOK_ENABLE_SECONDARY".to_string(), "on".to_string()),
251 (
252 "RUSTFS_NOTIFY_WEBHOOK_ENDPOINT_SECONDARY".to_string(),
253 "https://example.com/secondary".to_string(),
254 ),
255 ],
256 );
257
258 assert_eq!(instances.len(), 2);
259
260 let primary = instances
261 .iter()
262 .find(|instance| instance.instance_id == "primary")
263 .expect("primary notify instance should be normalized");
264 assert_eq!(primary.domain, TargetDomain::Notify);
265 assert_eq!(primary.plugin_id, "builtin:webhook");
266 assert!(primary.enabled);
267 assert_eq!(primary.effective_config.lookup(WEBHOOK_QUEUE_LIMIT).as_deref(), Some("42"));
268 assert_eq!(
269 primary.effective_config.lookup(WEBHOOK_ENDPOINT).as_deref(),
270 Some("https://example.com/primary")
271 );
272 assert_eq!(primary.source_hints.classification(), TargetInstanceSourceClass::Mixed);
273 assert!(primary.source_hints.has_file_default);
274 assert!(primary.source_hints.has_file_instance);
275 assert!(primary.source_hints.has_env_default);
276 assert!(!primary.source_hints.has_env_instance);
277
278 let secondary = instances
279 .iter()
280 .find(|instance| instance.instance_id == "secondary")
281 .expect("secondary env notify instance should be normalized");
282 assert!(secondary.enabled);
283 assert_eq!(
284 secondary.effective_config.lookup(WEBHOOK_ENDPOINT).as_deref(),
285 Some("https://example.com/secondary")
286 );
287 assert_eq!(secondary.effective_config.lookup(WEBHOOK_QUEUE_LIMIT).as_deref(), Some("42"));
288 assert_eq!(secondary.source_hints.classification(), TargetInstanceSourceClass::Mixed);
289 assert!(secondary.source_hints.has_file_default);
290 assert!(!secondary.source_hints.has_file_instance);
291 assert!(secondary.source_hints.has_env_default);
292 assert!(secondary.source_hints.has_env_instance);
293 }
294
295 #[test]
296 fn normalize_audit_instances_preserves_domain_and_subsystem() {
297 let mut cfg = Config(HashMap::new());
298 let mut subsystem = HashMap::new();
299
300 let mut default_kvs = KVS::new();
301 default_kvs.insert(ENABLE_KEY.to_string(), "off".to_string());
302 subsystem.insert("_".to_string(), default_kvs);
303
304 let mut primary = KVS::new();
305 primary.insert(ENABLE_KEY.to_string(), "on".to_string());
306 primary.insert(WEBHOOK_ENDPOINT.to_string(), "https://example.com/audit".to_string());
307 subsystem.insert("primary".to_string(), primary);
308
309 cfg.0.insert(AUDIT_WEBHOOK_SUB_SYS.to_string(), subsystem);
310
311 let instances = normalize_legacy_target_instances_from_env(&cfg, &audit_webhook_descriptor(), Vec::new());
312
313 assert_eq!(instances.len(), 1);
314 let primary = &instances[0];
315 assert_eq!(primary.domain, TargetDomain::Audit);
316 assert_eq!(primary.target_type, "webhook");
317 assert_eq!(primary.subsystem, AUDIT_WEBHOOK_SUB_SYS);
318 assert_eq!(primary.instance_id, "primary");
319 assert!(primary.enabled);
320 assert_eq!(
321 primary.effective_config.lookup(WEBHOOK_ENDPOINT).as_deref(),
322 Some("https://example.com/audit")
323 );
324 assert_eq!(primary.source_hints.classification(), TargetInstanceSourceClass::Config);
325 }
326
327 #[test]
328 fn normalize_instances_keeps_disabled_records() {
329 let cfg = Config(HashMap::new());
330
331 let instances = normalize_legacy_target_instances_from_env(
332 &cfg,
333 ¬ify_webhook_descriptor(),
334 vec![
335 ("RUSTFS_NOTIFY_WEBHOOK_ENABLE_DISABLED".to_string(), "off".to_string()),
336 (
337 "RUSTFS_NOTIFY_WEBHOOK_ENDPOINT_DISABLED".to_string(),
338 "https://example.com/disabled".to_string(),
339 ),
340 ],
341 );
342
343 assert_eq!(instances.len(), 1);
344 let disabled = &instances[0];
345 assert_eq!(disabled.instance_id, "disabled");
346 assert!(!disabled.enabled);
347 assert_eq!(disabled.source_hints.classification(), TargetInstanceSourceClass::Env);
348 assert!(disabled.source_hints.has_env_instance);
349 }
350
351 #[test]
352 fn normalize_instances_excludes_default_only_entries() {
353 let mut cfg = Config(HashMap::new());
354 let mut subsystem = HashMap::new();
355
356 let mut default_kvs = KVS::new();
357 default_kvs.insert(ENABLE_KEY.to_string(), "on".to_string());
358 default_kvs.insert(WEBHOOK_QUEUE_LIMIT.to_string(), "99".to_string());
359 subsystem.insert("_".to_string(), default_kvs);
360
361 cfg.0.insert(NOTIFY_WEBHOOK_SUB_SYS.to_string(), subsystem);
362
363 let instances = normalize_legacy_target_instances_from_env(
364 &cfg,
365 ¬ify_webhook_descriptor(),
366 vec![("RUSTFS_NOTIFY_WEBHOOK_QUEUE_LIMIT".to_string(), "100".to_string())],
367 );
368
369 assert!(instances.is_empty());
370 }
371
372 #[test]
373 fn compatibility_wrapper_matches_canonical_instance_model() {
374 let mut cfg = Config(HashMap::new());
375 let mut subsystem = HashMap::new();
376
377 let mut primary = KVS::new();
378 primary.insert(ENABLE_KEY.to_string(), "on".to_string());
379 primary.insert(WEBHOOK_ENDPOINT.to_string(), "https://example.com/primary".to_string());
380 subsystem.insert("primary".to_string(), primary);
381 cfg.0.insert(NOTIFY_WEBHOOK_SUB_SYS.to_string(), subsystem);
382
383 let descriptor = notify_webhook_descriptor();
384 let env = vec![("RUSTFS_NOTIFY_WEBHOOK_QUEUE_LIMIT".to_string(), "7".to_string())];
385
386 let canonical = try_normalize_target_plugin_instances_from_env(&cfg, &descriptor, env.clone())
387 .expect("canonical normalization should succeed");
388 let compatibility = normalize_legacy_target_instances_from_env(&cfg, &descriptor, env);
389
390 assert_eq!(canonical, compatibility);
391 }
392
393 #[test]
394 fn normalize_instances_rejects_invalid_enable_value() {
395 let error = try_normalize_target_plugin_instances_from_env(
396 &Config(HashMap::new()),
397 ¬ify_webhook_descriptor(),
398 vec![("RUSTFS_NOTIFY_WEBHOOK_ENABLE_PRIMARY".to_string(), "invalid".to_string())],
399 )
400 .expect_err("invalid enable value must be propagated by the public normalizer");
401
402 match error {
403 TargetError::Configuration(detail) => assert_eq!(detail, "Invalid enable value 'invalid'"),
404 other => panic!("expected a configuration error, got {other}"),
405 }
406 }
407
408 #[test]
409 fn legacy_normalizer_keeps_valid_instance_when_one_enable_is_invalid() {
410 let instances = normalize_legacy_target_instances_from_env(
411 &Config(HashMap::new()),
412 ¬ify_webhook_descriptor(),
413 vec![
414 ("RUSTFS_NOTIFY_WEBHOOK_ENABLE_GOOD".to_string(), "on".to_string()),
415 ("RUSTFS_NOTIFY_WEBHOOK_ENDPOINT_GOOD".to_string(), "https://example.com/good".to_string()),
416 ("RUSTFS_NOTIFY_WEBHOOK_ENABLE_BAD".to_string(), "invalid".to_string()),
417 ],
418 );
419
420 assert_eq!(instances.len(), 2);
421 let good = instances
422 .iter()
423 .find(|instance| instance.instance_id == "good")
424 .expect("valid instance should remain present");
425 let bad = instances
426 .iter()
427 .find(|instance| instance.instance_id == "bad")
428 .expect("invalid legacy instance should remain visible");
429 assert!(good.enabled);
430 assert!(!bad.enabled);
431 }
432}