1use crate::{AuditEntry, AuditError, AuditResult, factory::builtin_target_plugins};
16use rustfs_config::audit::AUDIT_ROUTE_PREFIX;
17use rustfs_config::server_config::{Config, KVS};
18use rustfs_targets::arn::TargetID;
19use rustfs_targets::{SharedTarget, Target, TargetError, TargetPluginRegistry, TargetRuntimeManager};
20use tracing::info;
21
22const LOG_COMPONENT_AUDIT: &str = "audit";
23const LOG_SUBSYSTEM_REGISTRY: &str = "registry";
24const EVENT_AUDIT_TARGET_REGISTRY_KEY_CREATED: &str = "audit_target_registry_key_created";
25const EVENT_AUDIT_TARGET_REGISTRY_STATE: &str = "audit_target_registry_state";
26
27pub struct AuditRegistry {
29 targets: TargetRuntimeManager<AuditEntry>,
31 plugins: TargetPluginRegistry<AuditEntry>,
33}
34
35impl Default for AuditRegistry {
36 fn default() -> Self {
37 Self::new()
38 }
39}
40
41impl AuditRegistry {
42 pub fn new() -> Self {
44 let mut plugins = TargetPluginRegistry::new();
45 plugins.register_all(builtin_target_plugins());
46
47 AuditRegistry {
48 targets: TargetRuntimeManager::new(),
49 plugins,
50 }
51 }
52
53 pub fn supports_target_type(&self, target_type: &str) -> bool {
54 self.plugins.supports_target_type(target_type)
55 }
56
57 pub async fn create_target(
67 &self,
68 target_type: &str,
69 id: String,
70 config: &KVS,
71 ) -> Result<Box<dyn Target<AuditEntry> + Send + Sync>, TargetError> {
72 self.plugins.create_target(target_type, id, config)
73 }
74
75 pub async fn create_audit_targets_from_config(
85 &self,
86 config: &Config,
87 ) -> AuditResult<Vec<Box<dyn Target<AuditEntry> + Send + Sync>>> {
88 self.plugins
89 .create_targets_from_config(config, AUDIT_ROUTE_PREFIX)
90 .await
91 .map_err(AuditError::from)
92 }
93
94 pub fn add_target(&mut self, _id: String, target: Box<dyn Target<AuditEntry> + Send + Sync>) {
100 debug_assert_eq!(_id, target.id().to_string());
101 self.targets.add_boxed(target);
102 }
103
104 pub fn add_shared_target(&mut self, _id: String, target: SharedTarget<AuditEntry>) {
105 debug_assert_eq!(_id, target.id().to_string());
106 self.targets.add_arc(target);
107 }
108
109 pub async fn remove_target(&mut self, id: &str) -> Option<SharedTarget<AuditEntry>> {
117 self.targets.remove_and_close(id).await
118 }
119
120 pub fn get_target(&self, id: &str) -> Option<SharedTarget<AuditEntry>> {
128 self.targets.get(id)
129 }
130
131 pub fn list_target_values(&self) -> Vec<SharedTarget<AuditEntry>> {
133 self.targets.values()
134 }
135
136 pub fn runtime_manager(&self) -> &TargetRuntimeManager<AuditEntry> {
137 &self.targets
138 }
139
140 pub fn runtime_manager_mut(&mut self) -> &mut TargetRuntimeManager<AuditEntry> {
141 &mut self.targets
142 }
143
144 pub fn list_targets(&self) -> Vec<String> {
149 self.targets.keys()
150 }
151
152 pub async fn close_all(&mut self) -> AuditResult<()> {
157 let mut first_error = None;
158
159 for target_id in self.targets.keys() {
160 if let Some(target) = self.targets.remove(&target_id)
161 && let Err(err) = target.close().await
162 {
163 tracing::error!(
164 event = EVENT_AUDIT_TARGET_REGISTRY_STATE,
165 component = LOG_COMPONENT_AUDIT,
166 subsystem = LOG_SUBSYSTEM_REGISTRY,
167 target_id = %target_id,
168 state = "close_failed",
169 error = %err,
170 "Failed to close target during shutdown"
171 );
172 if first_error.is_none() {
173 first_error = Some(err);
174 }
175 }
176 }
177
178 match first_error {
179 Some(err) => Err(AuditError::Target(err)),
180 None => Ok(()),
181 }
182 }
183
184 pub fn create_key(&self, target_type: &str, target_id: &str) -> String {
193 let key = TargetID::new(target_id.to_string(), target_type.to_string());
194 info!(
195 event = EVENT_AUDIT_TARGET_REGISTRY_KEY_CREATED,
196 component = LOG_COMPONENT_AUDIT,
197 subsystem = LOG_SUBSYSTEM_REGISTRY,
198 target_type = %target_type,
199 target_id = %target_id,
200 registry_key = %key,
201 "audit target registry state"
202 );
203 key.to_string()
204 }
205
206 pub fn enable_target(&self, target_type: &str, target_id: &str) -> AuditResult<()> {
215 let key = self.create_key(target_type, target_id);
216 if self.get_target(&key).is_some() {
217 info!(
218 event = EVENT_AUDIT_TARGET_REGISTRY_STATE,
219 component = LOG_COMPONENT_AUDIT,
220 subsystem = LOG_SUBSYSTEM_REGISTRY,
221 target_type = %target_type,
222 target_id = %target_id,
223 state = "enabled",
224 "audit target registry state"
225 );
226 Ok(())
227 } else {
228 Err(AuditError::Configuration(
229 format!("Target not found: {}-{}", target_type, target_id),
230 None,
231 ))
232 }
233 }
234
235 pub fn disable_target(&self, target_type: &str, target_id: &str) -> AuditResult<()> {
244 let key = self.create_key(target_type, target_id);
245 if self.get_target(&key).is_some() {
246 info!(
247 event = EVENT_AUDIT_TARGET_REGISTRY_STATE,
248 component = LOG_COMPONENT_AUDIT,
249 subsystem = LOG_SUBSYSTEM_REGISTRY,
250 target_type = %target_type,
251 target_id = %target_id,
252 state = "disabled",
253 "audit target registry state"
254 );
255 Ok(())
256 } else {
257 Err(AuditError::Configuration(
258 format!("Target not found: {}-{}", target_type, target_id),
259 None,
260 ))
261 }
262 }
263
264 pub fn upsert_target(
274 &mut self,
275 target_type: &str,
276 target_id: &str,
277 target: Box<dyn Target<AuditEntry> + Send + Sync>,
278 ) -> AuditResult<()> {
279 let key = self.create_key(target_type, target_id);
280 debug_assert_eq!(key, target.id().to_string());
281 self.targets.add_boxed(target);
282 Ok(())
283 }
284}
285
286#[cfg(test)]
287mod tests {
288 use super::AuditRegistry;
289 use crate::AuditError;
290 use rustfs_targets::TargetError;
291 use rustfs_targets::target::ChannelTargetType;
292 use rustfs_targets::testkit::MockTarget;
293
294 #[test]
295 fn registry_registers_amqp_factory() {
296 let registry = AuditRegistry::new();
297
298 assert!(registry.supports_target_type(ChannelTargetType::Amqp.as_str()));
299 }
300
301 #[tokio::test]
302 async fn close_all_returns_first_error_and_clears_targets() {
303 let mut registry = AuditRegistry::new();
304 let ok = MockTarget::new("ok", "webhook");
305 let ok_observer = ok.clone();
306 let fail = MockTarget::new("fail", "webhook")
307 .with_close_failures(usize::MAX)
308 .with_close_failure_error(|| TargetError::Unknown("close failed".to_string()));
309 let fail_observer = fail.clone();
310
311 registry.add_target(ok.target_id().to_string(), Box::new(ok));
312 registry.add_target(fail.target_id().to_string(), Box::new(fail));
313
314 let result = registry.close_all().await;
315
316 assert!(matches!(result, Err(AuditError::Target(TargetError::Unknown(_)))));
317 assert_eq!(ok_observer.close_call_count(), 1);
318 assert_eq!(fail_observer.close_call_count(), 1);
319 assert!(registry.list_targets().is_empty());
320 }
321}