Skip to main content

rustfs_targets/target/nats/
mod.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 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    // JetStream publish settings. Absent maps to off and the defaults.
73    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
193/// Returns true when the target sends credentials over a connection that does not require TLS, which
194/// would transmit the secrets in cleartext (backlog#983). TLS is active when tls_required is set or
195/// the address uses the tls:// scheme.
196fn 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    /// Cached JetStream context, one per target shared across clones so a single acker and semaphore
242    /// serve every clone. Read out under the lock before each publish await, never held across it.
243    /// Carries the stream-validation verdict for the bound connection.
244    jetstream_context: Arc<Mutex<Option<CachedJetStreamContext>>>,
245    tls_state: Arc<parking_lot::Mutex<TargetTlsState>>,
246    /// When set, the coordinator drives TLS reload, otherwise inline fingerprint change detection.
247    tls_adapter: Option<TlsReloadAdapter<async_nats::Client>>,
248    /// The concrete queue store, held typed so the target projects both the generic Store handle and
249    /// its own failed-events capability from a single store. Shared across clones through the Arc.
250    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        // Adapter-managed path: use the material directly from the TLS reload adapter.
316        if let Some(adapter) = &self.tls_adapter {
317            let client: async_nats::Client = (*adapter.current_material()).clone();
318
319            // Ensure the client is also stored locally so that close() can drain it.
320            {
321                let mut guard = self.client.lock().await;
322                *guard = Some(client.clone());
323            }
324            return Ok(client);
325        }
326
327        // Inline fingerprint fallback path (no coordinator).
328        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        // Swap in a context built from the winning client while the client lock is held, so client
356        // and context swap as a unit and no clone reads a context bound to the old client. Only the
357        // JetStream path caches a context. The previous context is drained after the locks release.
358        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        // Advance the recorded fingerprint only after the reconnect and flush succeed, so a rotation
370        // whose first reconnect fails stays detected and the next success still rebuilds the context.
371        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            // Mint the dedup id once at enqueue so it is identical across every retry and replay and
395            // distinct across entries. The body-hash fallback only covers entries queued before enable.
396            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            // publish only enqueues the message on the client's outbound channel. Flush to confirm the
410            // message reached the server before delivery is treated as successful (backlog#971).
411            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        // With JetStream enabled the health answer covers the stream as well as the connection, so a
440        // reachable broker with a failing stream reports the validation error. The verdict is cached
441        // and reset on a reconnect, TLS rotation, wrong-stream ack, or stream-not-found outcome, so a
442        // reset forces a live lookup on the next check.
443        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        // Drain the cached context's in-flight acks before dropping it so the async-nats acker exits.
489        // Drain the context before the client, since the acker rides the client connection.
490        let context = {
491            let mut guard = self.jetstream_context.lock().await;
492            guard.take()
493        };
494        // Bound the drain by the ack timeout so an unresponsive broker cannot stall close. On elapse
495        // the client drain below still runs.
496        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        // The durable queue store stays on disk, so a queued entry survives close and replays. Close
516        // releases the cached handles, never the entries.
517        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            // Record the verdict on the context so the first publish does not repeat the stream
573            // lookup that just passed.
574            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/// Coordinated TLS hot-reload implementation for NATS targets.
596///
597/// The coordinator calls these methods on a background poll loop to detect
598/// TLS file changes and rebuild the NATS client without restarting.
599#[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    // Absolute on Linux, macOS, and Windows. temp_dir needs no filesystem to exist for a
642    // validation-only test, and Path::is_absolute stays true across platforms.
643    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    /// A stream configuration that passes every assertion: it captures the configured subject, returns
719    /// acks, accepts writes, and deduplicates well beyond the retry lifetime.
720    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    // Live-broker behaviour tests. They require a NATS server with JetStream enabled and are ignored by
732    // default. To run locally:
733    //
734    //     docker run -d --name rustfs-nats-test -p 4222:4222 nats:2 -js
735    //     cargo test -p rustfs-targets --lib -- --ignored nats::tests::tls_change
736    //
737    // Override the server URL with RUSTFS_TEST_NATS_URL.
738    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    /// Log sink for asserting a specific line is written. The subscriber writes into a shared
743    /// buffer the assertion reads back.
744    #[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        // Token auth over a plaintext nats:// address without tls_required leaks the credential (backlog#983).
822        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        // tls_required protects the credentials.
830        let with_tls = NATSArgs {
831            tls_required: true,
832            ..insecure.clone()
833        };
834        assert!(!nats_sends_credentials_without_tls(&with_tls));
835
836        // A tls:// address also counts as protected.
837        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        // No credentials configured: nothing to leak.
844        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        // With the flag off the payload carries no dedup id, so its encoded meta matches a pre-feature entry.
851        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        // No JetStream context is built on the flag-off path.
859        assert!(disabled.jetstream_context.lock().await.is_none(), "no context is built with the flag off");
860    }
861}