Skip to main content

rustfs_targets/config/
instance.rs

1// Copyright 2024 RustFS Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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            &notify_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            &notify_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            &notify_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            &notify_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            &notify_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}