1use crate::plugin::PluginEvent;
16use crate::{
17 StoreError, Target,
18 arn::TargetID,
19 error::TargetError,
20 runtime::tls::{
21 ReloadableTargetTls, TargetTlsGeneration, TargetTlsInputSet, TlsReloadAdapter, config::ReloadApplyMode,
22 validate_tls_material,
23 },
24 store::{FailedEventStore, Key, QueueStore, Store},
25 target::{
26 ChannelTargetType, EntityTarget, QueuedPayload, QueuedPayloadMeta, TargetDeliveryCounters, TargetDeliverySnapshot,
27 TargetTlsState, TargetType, build_queued_payload_with_records, build_target_tls_fingerprint, is_connectivity_error,
28 open_target_queue_store_typed, persist_queued_payload_to_store, redacted_secret, with_delivery_deadline,
29 },
30};
31use async_trait::async_trait;
32use rustfs_config::{
33 NATS_CREDENTIALS_FILE, NATS_JETSTREAM_ACK_TIMEOUT_DEFAULT_SECS, NATS_TLS_CA, NATS_TLS_CLIENT_CERT, NATS_TLS_CLIENT_KEY,
34};
35use std::fmt;
36use std::path::{Path, PathBuf};
37use std::str::FromStr;
38use std::sync::Arc;
39use std::sync::atomic::{AtomicBool, Ordering};
40use std::time::Duration;
41use tokio::sync::Mutex;
42use tracing::{error, info, instrument, warn};
43use uuid::Uuid;
44
45mod jetstream;
46mod publish_error;
47mod validation;
48
49use jetstream::{CachedJetStreamContext, drain_jetstream_context};
50use publish_error::{classify_nats_flush_error, classify_nats_publish_error};
51
52pub(crate) use jetstream::resolve_dedup_id;
53pub(crate) use validation::{validate_jetstream_settings, validate_jetstream_stream};
54
55const NATS_CORE_DELIVERY_TIMEOUT: Duration = Duration::from_secs(30);
56
57#[derive(Clone)]
58pub struct NATSArgs {
59 pub enable: bool,
60 pub address: String,
61 pub subject: String,
62 pub username: String,
63 pub password: String,
64 pub token: String,
65 pub credentials_file: String,
66 pub tls_ca: String,
67 pub tls_client_cert: String,
68 pub tls_client_key: String,
69 pub tls_required: bool,
70 pub queue_dir: String,
71 pub queue_limit: u64,
72 pub jetstream_enable: Option<bool>,
74 pub jetstream_stream_name: Option<String>,
75 pub jetstream_ack_timeout_secs: Option<u64>,
76 pub target_type: TargetType,
77}
78
79impl fmt::Debug for NATSArgs {
80 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
81 f.debug_struct("NATSArgs")
82 .field("enable", &self.enable)
83 .field("address", &self.address)
84 .field("subject", &self.subject)
85 .field("username", &self.username)
86 .field("password", &redacted_secret(&self.password))
87 .field("token", &redacted_secret(&self.token))
88 .field("credentials_file", &redacted_secret(&self.credentials_file))
89 .field("tls_ca", &self.tls_ca)
90 .field("tls_client_cert", &self.tls_client_cert)
91 .field("tls_client_key", &redacted_secret(&self.tls_client_key))
92 .field("tls_required", &self.tls_required)
93 .field("queue_dir", &self.queue_dir)
94 .field("queue_limit", &self.queue_limit)
95 .field("jetstream_enable", &self.jetstream_enable)
96 .field("jetstream_stream_name", &self.jetstream_stream_name)
97 .field("jetstream_ack_timeout_secs", &self.jetstream_ack_timeout_secs)
98 .field("target_type", &self.target_type)
99 .finish()
100 }
101}
102
103impl NATSArgs {
104 pub fn validate(&self) -> Result<(), TargetError> {
105 if !self.enable {
106 return Ok(());
107 }
108
109 validate_nats_address(&self.address)?;
110 validate_nats_auth(self)?;
111
112 if self.subject.trim().is_empty() || self.subject.chars().any(char::is_whitespace) {
113 return Err(TargetError::Configuration(
114 "NATS subject cannot be empty or contain whitespace".to_string(),
115 ));
116 }
117
118 if !self.credentials_file.is_empty() && !Path::new(&self.credentials_file).is_absolute() {
119 return Err(TargetError::Configuration(format!("{NATS_CREDENTIALS_FILE} must be an absolute path")));
120 }
121 if !self.tls_ca.is_empty() && !Path::new(&self.tls_ca).is_absolute() {
122 return Err(TargetError::Configuration(format!("{NATS_TLS_CA} must be an absolute path")));
123 }
124 if !self.tls_client_cert.is_empty() && !Path::new(&self.tls_client_cert).is_absolute() {
125 return Err(TargetError::Configuration(format!("{NATS_TLS_CLIENT_CERT} must be an absolute path")));
126 }
127 if !self.tls_client_key.is_empty() && !Path::new(&self.tls_client_key).is_absolute() {
128 return Err(TargetError::Configuration(format!("{NATS_TLS_CLIENT_KEY} must be an absolute path")));
129 }
130 if self.tls_client_cert.is_empty() != self.tls_client_key.is_empty() {
131 return Err(TargetError::Configuration(
132 "NATS tls_client_cert and tls_client_key must be specified together".to_string(),
133 ));
134 }
135
136 if !self.queue_dir.is_empty() && !Path::new(&self.queue_dir).is_absolute() {
137 return Err(TargetError::Configuration("NATS queue directory must be an absolute path".to_string()));
138 }
139
140 validate_jetstream_settings(
141 self.jetstream_enable.unwrap_or(false),
142 self.jetstream_stream_name.as_deref().unwrap_or_default(),
143 &self.queue_dir,
144 self.jetstream_ack_timeout_secs,
145 )?;
146
147 Ok(())
148 }
149}
150
151pub fn validate_nats_address(address: &str) -> Result<async_nats::ServerAddr, TargetError> {
152 let server = async_nats::ServerAddr::from_str(address)
153 .map_err(|e| TargetError::Configuration(format!("Invalid NATS address: {e}")))?;
154
155 if server.has_user_pass() {
156 return Err(TargetError::Configuration("NATS address must not embed username or password".to_string()));
157 }
158
159 Ok(server)
160}
161
162fn validate_nats_auth(args: &NATSArgs) -> Result<(), TargetError> {
163 let mut auth_methods = 0usize;
164
165 if !args.token.is_empty() {
166 auth_methods += 1;
167 }
168
169 if !args.credentials_file.is_empty() {
170 auth_methods += 1;
171 }
172
173 let has_user = !args.username.is_empty();
174 let has_password = !args.password.is_empty();
175 if has_user || has_password {
176 if has_user != has_password {
177 return Err(TargetError::Configuration(
178 "NATS username and password must be specified together".to_string(),
179 ));
180 }
181 auth_methods += 1;
182 }
183
184 if auth_methods > 1 {
185 return Err(TargetError::Configuration(
186 "NATS supports only one auth method at a time: token, username/password, or credentials_file".to_string(),
187 ));
188 }
189
190 Ok(())
191}
192
193fn nats_sends_credentials_without_tls(args: &NATSArgs) -> bool {
197 let has_auth = !args.token.is_empty() || !args.credentials_file.is_empty() || !args.username.is_empty();
198 if !has_auth {
199 return false;
200 }
201 let scheme_is_tls = args.address.trim_start().to_ascii_lowercase().starts_with("tls://");
202 !(args.tls_required || scheme_is_tls)
203}
204
205pub async fn connect_nats(args: &NATSArgs) -> Result<async_nats::Client, TargetError> {
206 args.validate()?;
207
208 let mut options = async_nats::ConnectOptions::new().require_tls(args.tls_required);
209
210 if !args.token.is_empty() {
211 options = options.token(args.token.clone());
212 } else if !args.username.is_empty() {
213 options = options.user_and_password(args.username.clone(), args.password.clone());
214 } else if !args.credentials_file.is_empty() {
215 options = options
216 .credentials_file(&args.credentials_file)
217 .await
218 .map_err(|e| TargetError::Configuration(format!("Failed to load NATS credentials file: {e}")))?;
219 }
220
221 if !args.tls_ca.is_empty() {
222 options = options.add_root_certificates(PathBuf::from(&args.tls_ca));
223 }
224 if !args.tls_client_cert.is_empty() {
225 options = options.add_client_certificate(PathBuf::from(&args.tls_client_cert), PathBuf::from(&args.tls_client_key));
226 }
227
228 options
229 .connect(args.address.clone())
230 .await
231 .map_err(|e| TargetError::Network(format!("Failed to connect to NATS server: {e}")))
232}
233
234pub struct NATSTarget<E>
235where
236 E: PluginEvent,
237{
238 id: TargetID,
239 args: NATSArgs,
240 client: Arc<Mutex<Option<async_nats::Client>>>,
241 jetstream_context: Arc<Mutex<Option<CachedJetStreamContext>>>,
245 tls_state: Arc<parking_lot::Mutex<TargetTlsState>>,
246 tls_adapter: Option<TlsReloadAdapter<async_nats::Client>>,
248 store: Option<Arc<QueueStore<QueuedPayload>>>,
251 connected: AtomicBool,
252 delivery_counters: Arc<TargetDeliveryCounters>,
253 _phantom: std::marker::PhantomData<E>,
254}
255
256impl<E> NATSTarget<E>
257where
258 E: PluginEvent,
259{
260 pub fn clone_box(&self) -> Box<dyn Target<E> + Send + Sync> {
261 Box::new(NATSTarget::<E> {
262 id: self.id.clone(),
263 args: self.args.clone(),
264 client: Arc::clone(&self.client),
265 jetstream_context: Arc::clone(&self.jetstream_context),
266 tls_state: Arc::clone(&self.tls_state),
267 tls_adapter: self.tls_adapter.clone(),
268 store: self.store.clone(),
269 connected: AtomicBool::new(self.connected.load(Ordering::SeqCst)),
270 delivery_counters: Arc::clone(&self.delivery_counters),
271 _phantom: std::marker::PhantomData,
272 })
273 }
274
275 #[instrument(skip(args), fields(target_id_as_string = %id))]
276 pub fn new(id: String, args: NATSArgs) -> Result<Self, TargetError> {
277 args.validate()?;
278 if args.enable && nats_sends_credentials_without_tls(&args) {
279 warn!(
280 target_id = %id,
281 address = %args.address,
282 "NATS target sends authentication credentials without TLS; secrets are transmitted in cleartext. Enable tls_required or use a tls:// address."
283 );
284 }
285 let target_id = TargetID::new(id, ChannelTargetType::Nats.as_str().to_string());
286 let queue_store = open_target_queue_store_typed(
287 &args.queue_dir,
288 args.queue_limit,
289 args.target_type,
290 ChannelTargetType::Nats.as_str(),
291 &target_id,
292 "Failed to open store for NATS target",
293 )?
294 .map(Arc::new);
295
296 Ok(Self {
297 id: target_id,
298 args,
299 client: Arc::new(Mutex::new(None)),
300 jetstream_context: Arc::new(Mutex::new(None)),
301 tls_state: Arc::new(parking_lot::Mutex::new(TargetTlsState::default())),
302 tls_adapter: None,
303 store: queue_store,
304 connected: AtomicBool::new(false),
305 delivery_counters: Arc::new(TargetDeliveryCounters::default()),
306 _phantom: std::marker::PhantomData,
307 })
308 }
309
310 async fn invalidate_cached_client_connection(&self) {
311 *self.client.lock().await = None;
312 }
313
314 async fn get_or_connect(&self) -> Result<async_nats::Client, TargetError> {
315 if let Some(adapter) = &self.tls_adapter {
317 let client: async_nats::Client = (*adapter.current_material()).clone();
318
319 {
321 let mut guard = self.client.lock().await;
322 *guard = Some(client.clone());
323 }
324 return Ok(client);
325 }
326
327 let next_fingerprint =
329 build_target_tls_fingerprint(&self.args.tls_ca, &self.args.tls_client_cert, &self.args.tls_client_key).await?;
330 let tls_changed = {
331 let tls_state_guard = self.tls_state.lock();
332 tls_state_guard.needs_update(&next_fingerprint)
333 };
334 if tls_changed {
335 self.invalidate_cached_client_connection().await;
336 }
337
338 {
339 let guard = self.client.lock().await;
340 if let Some(client) = guard.as_ref() {
341 return Ok(client.clone());
342 }
343 }
344
345 let client = connect_nats(&self.args).await?;
346 client
347 .flush()
348 .await
349 .map_err(|e| TargetError::Network(format!("Failed to flush NATS connection: {e}")))?;
350 self.connected.store(true, Ordering::SeqCst);
351
352 let mut client_guard = self.client.lock().await;
353 let shared = client_guard.get_or_insert_with(|| client.clone()).clone();
354
355 let previous_context = if tls_changed && self.jetstream_enabled() {
359 let ack_timeout = self.ack_timeout();
360 let mut context_guard = self.jetstream_context.lock().await;
361 let mut context = async_nats::jetstream::new(shared.clone());
362 context.set_timeout(ack_timeout);
363 context_guard.replace(CachedJetStreamContext::new(context))
364 } else {
365 None
366 };
367 drop(client_guard);
368
369 if tls_changed {
372 self.tls_state.lock().refresh(next_fingerprint);
373 }
374
375 drain_jetstream_context(previous_context.map(|cached| cached.context)).await;
376 Ok(shared)
377 }
378
379 fn jetstream_enabled(&self) -> bool {
380 self.args.jetstream_enable.unwrap_or(false)
381 }
382
383 fn ack_timeout(&self) -> Duration {
384 Duration::from_secs(
385 self.args
386 .jetstream_ack_timeout_secs
387 .unwrap_or(NATS_JETSTREAM_ACK_TIMEOUT_DEFAULT_SECS),
388 )
389 }
390
391 fn build_queued_payload(&self, event: &EntityTarget<E>) -> Result<QueuedPayload, TargetError> {
392 let mut queued = build_queued_payload_with_records(event, vec![event.clone()])?;
393 if self.jetstream_enabled() {
394 queued.meta.dedup_id = Uuid::new_v4().to_string();
397 }
398 Ok(queued)
399 }
400
401 async fn send_body(&self, body: Vec<u8>) -> Result<(), TargetError> {
402 let result = with_delivery_deadline(NATS_CORE_DELIVERY_TIMEOUT, "NATS delivery", async {
403 let client = self.get_or_connect().await?;
404 client
405 .publish(self.args.subject.clone(), body.into())
406 .await
407 .map_err(|err| classify_nats_publish_error(&err))?;
408
409 client.flush().await.map_err(|err| classify_nats_flush_error(&err))?;
412 Ok(())
413 })
414 .await;
415
416 if let Err(err) = result {
417 if is_connectivity_error(&err) {
418 self.invalidate_cached_client_connection().await;
419 self.connected.store(false, Ordering::SeqCst);
420 }
421 return Err(err);
422 }
423
424 self.delivery_counters.record_success();
425 Ok(())
426 }
427}
428
429#[async_trait]
430impl<E> Target<E> for NATSTarget<E>
431where
432 E: PluginEvent,
433{
434 fn id(&self) -> TargetID {
435 self.id.clone()
436 }
437
438 async fn is_active(&self) -> Result<bool, TargetError> {
439 if self.jetstream_enabled() {
444 self.validated_jetstream_context().await?;
445 }
446 let client = self.get_or_connect().await?;
447 client
448 .flush()
449 .await
450 .map_err(|e| TargetError::Network(format!("NATS health check failed: {e}")))?;
451 Ok(true)
452 }
453
454 async fn save(&self, event: Arc<EntityTarget<E>>) -> Result<(), TargetError> {
455 let queued = match self.build_queued_payload(&event) {
456 Ok(queued) => queued,
457 Err(err) => {
458 self.delivery_counters.record_final_failure();
459 return Err(err);
460 }
461 };
462
463 if let Some(store) = &self.store {
464 if let Err(e) = persist_queued_payload_to_store(store.as_ref(), &queued) {
465 self.delivery_counters.record_final_failure();
466 return Err(e);
467 }
468 Ok(())
469 } else {
470 if let Err(err) = self.send_body(queued.body).await {
471 self.delivery_counters.record_final_failure();
472 return Err(err);
473 }
474 Ok(())
475 }
476 }
477
478 async fn send_raw_from_store(&self, key: Key, body: Vec<u8>, meta: QueuedPayloadMeta) -> Result<(), TargetError> {
479 if self.jetstream_enabled() {
480 let dedup_id = resolve_dedup_id(&meta.dedup_id, &key);
481 self.publish_jetstream(body, &dedup_id).await
482 } else {
483 self.send_body(body).await
484 }
485 }
486
487 async fn close(&self) -> Result<(), TargetError> {
488 let context = {
491 let mut guard = self.jetstream_context.lock().await;
492 guard.take()
493 };
494 if tokio::time::timeout(self.ack_timeout(), drain_jetstream_context(context.map(|cached| cached.context)))
497 .await
498 .is_err()
499 {
500 warn!(target_id = %self.id, "Timed out draining JetStream acks on close, proceeding to close the client");
501 }
502
503 let client = {
504 let mut guard = self.client.lock().await;
505 guard.take()
506 };
507 self.tls_state.lock().reset();
508 self.connected.store(false, Ordering::SeqCst);
509 if let Some(client) = client {
510 client
511 .drain()
512 .await
513 .map_err(|e| TargetError::Network(format!("Failed to drain NATS client: {e}")))?;
514 }
515 info!(target_id = %self.id, "NATS target closed");
518 Ok(())
519 }
520
521 fn store(&self) -> Option<&(dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send + Sync)> {
522 self.store
523 .as_deref()
524 .map(|store| store as &(dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send + Sync))
525 }
526
527 fn failed_store(&self) -> Option<&dyn FailedEventStore> {
528 self.store.as_deref().map(|store| store as &dyn FailedEventStore)
529 }
530
531 async fn handle_terminal_failure(
532 &self,
533 store: &(dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send),
534 key: &Key,
535 error: &TargetError,
536 retry_count: u32,
537 ) -> bool {
538 let Some(failed_store) = self.failed_store() else {
539 error!(
540 target_id = %self.id,
541 replay_key = %key,
542 "NATS target has no failed-events store for the terminal move"
543 );
544 return false;
545 };
546 match jetstream::move_entry_to_failed_store(store, failed_store, &self.id, key, error, retry_count).await {
547 Ok(()) => true,
548 Err(move_err) => {
549 error!(
550 target_id = %self.id,
551 error = %move_err,
552 replay_key = %key,
553 "Failed to move event to the failed-events store"
554 );
555 false
556 }
557 }
558 }
559
560 fn clone_dyn(&self) -> Box<dyn Target<E> + Send + Sync> {
561 self.clone_box()
562 }
563
564 async fn init(&self) -> Result<(), TargetError> {
565 if !self.is_enabled() {
566 return Ok(());
567 }
568 let _ = self.get_or_connect().await?;
569 if self.jetstream_enabled() {
570 let cached = self.jetstream_context().await?;
571 validate_jetstream_stream(&cached.context, &self.args, &self.id.to_string(), Some(&cached.validation_logged)).await?;
572 cached.stream_validated.store(true, Ordering::Release);
575 }
576 Ok(())
577 }
578
579 fn is_enabled(&self) -> bool {
580 self.args.enable
581 }
582
583 fn delivery_snapshot(&self) -> TargetDeliverySnapshot {
584 self.delivery_counters.snapshot(
585 self.store.as_deref().map_or(0, |store| store.len() as u64),
586 self.failed_store().map_or(0, |failed_store| failed_store.failed_len() as u64),
587 )
588 }
589
590 fn record_final_failure(&self) {
591 self.delivery_counters.record_final_failure();
592 }
593}
594
595#[async_trait]
600impl<E> ReloadableTargetTls for NATSTarget<E>
601where
602 E: PluginEvent,
603{
604 type Material = async_nats::Client;
605
606 fn tls_input_set(&self) -> TargetTlsInputSet {
607 TargetTlsInputSet {
608 ca_path: self.args.tls_ca.clone(),
609 client_cert_path: self.args.tls_client_cert.clone(),
610 client_key_path: self.args.tls_client_key.clone(),
611 target_label: format!("nats:{}", self.id.id),
612 }
613 }
614
615 async fn build_tls_material(&self) -> Result<Self::Material, TargetError> {
616 connect_nats(&self.args).await
617 }
618
619 async fn apply_tls_material(
620 &self,
621 _generation: TargetTlsGeneration,
622 material: Arc<Self::Material>,
623 _mode: ReloadApplyMode,
624 ) -> Result<(), TargetError> {
625 let mut guard = self.client.lock().await;
626 *guard = Some((*material).clone());
627 Ok(())
628 }
629
630 async fn validate_tls_files(&self) -> Result<(), TargetError> {
631 validate_tls_material(&self.args.tls_ca, &self.args.tls_client_cert, &self.args.tls_client_key)
632 }
633}
634
635#[cfg(test)]
636pub(crate) mod test_support {
637 use super::*;
638 use async_nats::jetstream;
639 use rustfs_s3_types::EventName;
640
641 pub(crate) fn nats_queue_dir() -> String {
644 std::env::temp_dir().join("rustfs-nats-queue").to_string_lossy().into_owned()
645 }
646
647 pub(crate) fn base_args() -> NATSArgs {
648 NATSArgs {
649 enable: true,
650 address: "nats://127.0.0.1:4222".to_string(),
651 subject: "rustfs.events".to_string(),
652 username: String::new(),
653 password: String::new(),
654 token: String::new(),
655 credentials_file: String::new(),
656 tls_ca: String::new(),
657 tls_client_cert: String::new(),
658 tls_client_key: String::new(),
659 tls_required: false,
660 queue_dir: String::new(),
661 queue_limit: 0,
662 jetstream_enable: None,
663 jetstream_stream_name: None,
664 jetstream_ack_timeout_secs: None,
665 target_type: TargetType::NotifyEvent,
666 }
667 }
668
669 pub(crate) fn jetstream_args(target_type: TargetType) -> NATSArgs {
670 NATSArgs {
671 jetstream_enable: Some(true),
672 jetstream_stream_name: Some("RUSTFS_EVENTS".to_string()),
673 jetstream_ack_timeout_secs: Some(30),
674 target_type,
675 ..base_args()
676 }
677 }
678
679 pub(crate) fn nats_target(args: NATSArgs) -> NATSTarget<String> {
680 NATSTarget::<String> {
681 id: TargetID::new("test-target".to_string(), ChannelTargetType::Nats.as_str().to_string()),
682 args,
683 client: Arc::new(Mutex::new(None)),
684 jetstream_context: Arc::new(Mutex::new(None)),
685 tls_state: Arc::new(parking_lot::Mutex::new(TargetTlsState::default())),
686 tls_adapter: None,
687 store: None,
688 connected: AtomicBool::new(false),
689 delivery_counters: Arc::new(TargetDeliveryCounters::default()),
690 _phantom: std::marker::PhantomData,
691 }
692 }
693
694 pub(crate) fn sample_event() -> EntityTarget<String> {
695 EntityTarget {
696 object_name: "folder/object.txt".to_string(),
697 bucket_name: "bucket-a".to_string(),
698 event_name: EventName::ObjectCreatedPut,
699 data: "payload-data".to_string(),
700 }
701 }
702
703 pub(crate) fn sample_stored_key(name: &str) -> Key {
704 Key {
705 name: name.to_string(),
706 extension: ".event".to_string(),
707 item_count: 1,
708 compress: false,
709 }
710 }
711
712 pub(crate) fn nats_target_with_store(args: NATSArgs, queue_dir: &str) -> NATSTarget<String> {
713 let mut configured = args;
714 configured.queue_dir = queue_dir.to_string();
715 NATSTarget::<String>::new("store-target".to_string(), configured).expect("target with store builds")
716 }
717
718 pub(crate) fn writable_stream_config(subject: &str) -> jetstream::stream::Config {
721 jetstream::stream::Config {
722 name: "RUSTFS_EVENTS".to_string(),
723 subjects: vec![subject.to_string()],
724 no_ack: false,
725 sealed: false,
726 duplicate_window: Duration::from_secs(600),
727 ..Default::default()
728 }
729 }
730
731 pub(crate) fn broker_url() -> String {
739 std::env::var("RUSTFS_TEST_NATS_URL").unwrap_or_else(|_| "nats://127.0.0.1:4222".to_string())
740 }
741
742 #[derive(Clone, Default)]
745 pub(crate) struct CapturedLog(Arc<parking_lot::Mutex<Vec<u8>>>);
746
747 impl CapturedLog {
748 pub(crate) fn contents(&self) -> String {
749 String::from_utf8_lossy(&self.0.lock()).into_owned()
750 }
751 }
752
753 impl std::io::Write for CapturedLog {
754 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
755 self.0.lock().extend_from_slice(buf);
756 Ok(buf.len())
757 }
758
759 fn flush(&mut self) -> std::io::Result<()> {
760 Ok(())
761 }
762 }
763
764 impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for CapturedLog {
765 type Writer = CapturedLog;
766
767 fn make_writer(&'a self) -> Self::Writer {
768 self.clone()
769 }
770 }
771}
772
773#[cfg(test)]
774mod tests {
775 use super::*;
776 use crate::target::REDACTED_SECRET;
777 use crate::target::nats::test_support::*;
778
779 #[test]
780 fn debug_redacts_nats_secret_fields() {
781 let args = NATSArgs {
782 password: "nats-password".to_string(),
783 token: "nats-token".to_string(),
784 credentials_file: "/etc/rustfs/nats.creds".to_string(),
785 tls_client_key: "/etc/rustfs/nats.key".to_string(),
786 ..base_args()
787 };
788
789 let rendered = format!("{args:?}");
790
791 assert!(!rendered.contains("nats-password"));
792 assert!(!rendered.contains("nats-token"));
793 assert!(!rendered.contains("/etc/rustfs/nats.creds"));
794 assert!(!rendered.contains("/etc/rustfs/nats.key"));
795 assert!(rendered.contains(REDACTED_SECRET));
796 assert!(rendered.contains("rustfs.events"));
797 }
798
799 #[test]
800 fn validate_nats_rejects_multiple_auth_methods() {
801 let args = NATSArgs {
802 token: "abc".to_string(),
803 username: "user".to_string(),
804 password: "pass".to_string(),
805 ..base_args()
806 };
807 assert!(args.validate().is_err());
808 }
809
810 #[test]
811 fn validate_nats_rejects_relative_queue_dir() {
812 let args = NATSArgs {
813 queue_dir: "relative/path".to_string(),
814 ..base_args()
815 };
816 assert!(args.validate().is_err());
817 }
818
819 #[test]
820 fn nats_credentials_without_tls_is_detected() {
821 let insecure = NATSArgs {
823 token: "secret-token".to_string(),
824 tls_required: false,
825 ..base_args()
826 };
827 assert!(nats_sends_credentials_without_tls(&insecure));
828
829 let with_tls = NATSArgs {
831 tls_required: true,
832 ..insecure.clone()
833 };
834 assert!(!nats_sends_credentials_without_tls(&with_tls));
835
836 let tls_scheme = NATSArgs {
838 address: "tls://127.0.0.1:4222".to_string(),
839 ..insecure
840 };
841 assert!(!nats_sends_credentials_without_tls(&tls_scheme));
842
843 let no_auth = base_args();
845 assert!(!nats_sends_credentials_without_tls(&no_auth));
846 }
847
848 #[tokio::test]
849 async fn flag_off_stored_meta_is_byte_identical_to_pre_feature() {
850 let disabled = nats_target(base_args());
852 let off_payload = disabled.build_queued_payload(&sample_event()).expect("payload builds");
853 assert!(off_payload.meta.dedup_id.is_empty(), "no dedup id is minted with the flag off");
854
855 let meta_json = serde_json::to_string(&off_payload.meta).expect("meta serializes");
856 assert!(!meta_json.contains("dedup_id"), "the absent dedup id is skipped on serialization");
857
858 assert!(disabled.jetstream_context.lock().await.is_none(), "no context is built with the flag off");
860 }
861}