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::{Key, Store},
25 target::{
26 ChannelTargetType, EntityTarget, QueuedPayload, QueuedPayloadMeta, TargetDeliveryCounters, TargetDeliverySnapshot,
27 TargetType, build_queued_payload, invalidate_cache_on_connectivity_error, is_connectivity_error,
28 mark_target_disconnected_on_connectivity_error, open_target_queue_store, persist_queued_payload_to_store,
29 with_delivery_deadline,
30 },
31};
32use async_trait::async_trait;
33use redis::{
34 AsyncCommands, Client, ClientTlsConfig, ConnectionInfo, IntoConnectionInfo, RedisError, TlsCertificates,
35 aio::{ConnectionManager, ConnectionManagerConfig},
36 cmd,
37 io::tcp::{TcpSettings, socket2},
38};
39use rustfs_config::{REDIS_TLS_CA, REDIS_TLS_CLIENT_CERT, REDIS_TLS_CLIENT_KEY, REDIS_TLS_POLICY};
40use rustls::pki_types::CertificateDer;
41use rustls::pki_types::pem::PemObject;
42use std::fmt;
43use std::io::BufReader;
44use std::path::Path;
45use std::sync::Arc;
46use std::sync::atomic::{AtomicBool, Ordering};
47use std::time::Duration;
48use tokio::sync::Mutex;
49use tracing::{debug, info, instrument, warn};
50use url::Url;
51
52const REDIS_CONNECTION_TIMEOUT_DEFAULT: Duration = Duration::from_secs(5);
53const REDIS_RESPONSE_TIMEOUT_DEFAULT: Duration = Duration::from_secs(5);
54
55fn redis_total_delivery_timeout(args: &RedisArgs) -> Duration {
56 let attempts = u32::try_from(args.max_retry_attempts).unwrap_or(u32::MAX);
57 let per_attempt = args
58 .connection_timeout
59 .unwrap_or(REDIS_CONNECTION_TIMEOUT_DEFAULT)
60 .saturating_add(args.response_timeout.unwrap_or(REDIS_RESPONSE_TIMEOUT_DEFAULT))
61 .saturating_add(args.max_retry_delay.unwrap_or(Duration::from_secs(2)));
62 per_attempt.saturating_mul(attempts)
63}
64
65#[derive(Debug, Clone, Copy, PartialEq, Eq)]
66pub enum RedisTlsPolicy {
67 SystemCa,
68 CustomCa,
69}
70
71impl RedisTlsPolicy {
72 fn parse(value: &str) -> Result<Self, TargetError> {
73 match value.trim() {
74 value if value.eq_ignore_ascii_case("system_ca") => Ok(Self::SystemCa),
75 value if value.eq_ignore_ascii_case("custom_ca") => Ok(Self::CustomCa),
76 _ => Err(TargetError::Configuration(
77 "Redis tls_policy must be one of: system_ca, custom_ca".to_string(),
78 )),
79 }
80 }
81}
82
83#[derive(Debug, Clone, Default, PartialEq, Eq)]
84pub struct RedisTlsConfig {
85 pub policy: Option<RedisTlsPolicy>,
86 pub ca_path: String,
87 pub client_cert_path: String,
88 pub client_key_path: String,
89 pub allow_insecure: bool,
90}
91
92impl RedisTlsConfig {
93 pub fn from_values(
94 policy: Option<&str>,
95 ca_path: Option<&str>,
96 client_cert_path: Option<&str>,
97 client_key_path: Option<&str>,
98 allow_insecure: Option<&str>,
99 ) -> Result<Self, TargetError> {
100 let policy = match policy.map(str::trim).filter(|value| !value.is_empty()) {
101 Some(value) => Some(RedisTlsPolicy::parse(value)?),
102 None => None,
103 };
104 let allow_insecure = allow_insecure
105 .map(str::trim)
106 .filter(|value| !value.is_empty())
107 .map(|value| {
108 value
109 .parse::<rustfs_config::EnableState>()
110 .map(rustfs_config::EnableState::is_enabled)
111 .or_else(|_| value.parse::<bool>())
112 .map_err(|_| TargetError::Configuration("Redis tls_allow_insecure must be a boolean value".to_string()))
113 })
114 .transpose()?
115 .unwrap_or(false);
116
117 Ok(Self {
118 policy,
119 ca_path: ca_path.unwrap_or_default().trim().to_string(),
120 client_cert_path: client_cert_path.unwrap_or_default().trim().to_string(),
121 client_key_path: client_key_path.unwrap_or_default().trim().to_string(),
122 allow_insecure,
123 })
124 }
125}
126
127#[derive(Clone)]
128pub struct RedisArgs {
129 pub enable: bool,
131 pub url: Url,
133 pub channel: String,
135 pub username: Option<String>,
137 pub password: Option<String>,
139 pub tls: RedisTlsConfig,
141 pub keep_alive: Duration,
143 pub queue_dir: String,
145 pub queue_limit: u64,
147 pub max_retry_attempts: usize,
149 pub reconnect_retry_attempts: Option<usize>,
151 pub min_retry_delay: Option<Duration>,
153 pub max_retry_delay: Option<Duration>,
155 pub connection_timeout: Option<Duration>,
157 pub response_timeout: Option<Duration>,
159 pub pipeline_buffer_size: Option<usize>,
161 pub target_type: TargetType,
163}
164
165fn redact_redis_url(url: &Url) -> String {
166 let mut redacted = url.clone();
167 if redacted.password().is_some() {
168 let _ = redacted.set_password(Some("***"));
169 }
170 redacted.to_string()
171}
172
173impl fmt::Debug for RedisArgs {
174 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
175 f.debug_struct("RedisArgs")
176 .field("enable", &self.enable)
177 .field("url", &redact_redis_url(&self.url))
178 .field("channel", &self.channel)
179 .field("username", &self.username)
180 .field(
181 "password",
182 if self.password.as_deref().unwrap_or_default().is_empty() {
183 &""
184 } else {
185 &"***REDACTED***"
186 },
187 )
188 .field("tls", &self.tls)
189 .field("keep_alive", &self.keep_alive)
190 .field("queue_dir", &self.queue_dir)
191 .field("queue_limit", &self.queue_limit)
192 .field("max_retry_attempts", &self.max_retry_attempts)
193 .field("reconnect_retry_attempts", &self.reconnect_retry_attempts)
194 .field("min_retry_delay", &self.min_retry_delay)
195 .field("max_retry_delay", &self.max_retry_delay)
196 .field("connection_timeout", &self.connection_timeout)
197 .field("response_timeout", &self.response_timeout)
198 .field("pipeline_buffer_size", &self.pipeline_buffer_size)
199 .field("target_type", &self.target_type)
200 .finish()
201 }
202}
203
204impl RedisArgs {
205 pub fn validate(&self) -> Result<(), TargetError> {
206 if !self.enable {
207 return Ok(());
208 }
209
210 validate_redis_url(&self.url)?;
211 validate_redis_tls_config(&self.url, &self.tls)?;
212
213 if self.channel.trim().is_empty() {
214 return Err(TargetError::Configuration("Redis channel cannot be empty".to_string()));
215 }
216
217 if self.username.as_deref().unwrap_or_default().is_empty() != self.password.as_deref().unwrap_or_default().is_empty()
218 && !(self.username.is_none() && self.password.is_none())
219 {
220 return Err(TargetError::Configuration(
221 "Redis username and password must be specified together when provided explicitly".to_string(),
222 ));
223 }
224
225 if self.max_retry_attempts == 0 {
226 return Err(TargetError::Configuration(
227 "Redis max_retry_attempts must be greater than zero".to_string(),
228 ));
229 }
230
231 if self.connection_timeout == Some(Duration::ZERO) {
232 return Err(TargetError::Configuration(
233 "Redis connection_timeout must be greater than zero".to_string(),
234 ));
235 }
236
237 if self.response_timeout == Some(Duration::ZERO) {
238 return Err(TargetError::Configuration("Redis response_timeout must be greater than zero".to_string()));
239 }
240
241 if self.pipeline_buffer_size == Some(0) {
242 return Err(TargetError::Configuration(
243 "Redis pipeline_buffer_size must be greater than zero".to_string(),
244 ));
245 }
246
247 if let (Some(min_retry_delay), Some(max_retry_delay)) = (self.min_retry_delay, self.max_retry_delay)
248 && max_retry_delay < min_retry_delay
249 {
250 return Err(TargetError::Configuration(
251 "Redis max_retry_delay must be greater than or equal to min_retry_delay".to_string(),
252 ));
253 }
254
255 if !self.queue_dir.is_empty() && !Path::new(&self.queue_dir).is_absolute() {
256 return Err(TargetError::Configuration("Redis queue_dir path should be absolute".to_string()));
257 }
258
259 Ok(())
260 }
261}
262
263pub fn validate_redis_url(url: &Url) -> Result<(), TargetError> {
264 let _: ConnectionInfo = url.clone().into_connection_info().map_err(map_redis_error)?;
265 Ok(())
266}
267
268fn validate_redis_tls_config(url: &Url, tls: &RedisTlsConfig) -> Result<(), TargetError> {
269 let secure_scheme = matches!(url.scheme(), "rediss" | "valkeys");
270
271 if !tls.client_cert_path.is_empty() && !Path::new(&tls.client_cert_path).is_absolute() {
272 return Err(TargetError::Configuration(format!("{REDIS_TLS_CLIENT_CERT} must be an absolute path")));
273 }
274 if !tls.client_key_path.is_empty() && !Path::new(&tls.client_key_path).is_absolute() {
275 return Err(TargetError::Configuration(format!("{REDIS_TLS_CLIENT_KEY} must be an absolute path")));
276 }
277 if tls.client_cert_path.is_empty() != tls.client_key_path.is_empty() {
278 return Err(TargetError::Configuration(
279 "Redis tls_client_cert and tls_client_key must be specified together".to_string(),
280 ));
281 }
282
283 if !secure_scheme {
284 if tls.policy.is_some()
285 || !tls.ca_path.is_empty()
286 || !tls.client_cert_path.is_empty()
287 || !tls.client_key_path.is_empty()
288 || tls.allow_insecure
289 {
290 return Err(TargetError::Configuration(
291 "TLS settings are only allowed for rediss/valkeys schemes".to_string(),
292 ));
293 }
294 return Ok(());
295 }
296
297 if let Some(policy) = tls.policy {
298 match policy {
299 RedisTlsPolicy::SystemCa => {
300 if !tls.ca_path.is_empty() {
301 return Err(TargetError::Configuration(format!(
302 "{REDIS_TLS_CA} is not allowed when {REDIS_TLS_POLICY}=system_ca"
303 )));
304 }
305 }
306 RedisTlsPolicy::CustomCa => {
307 if tls.ca_path.is_empty() {
308 return Err(TargetError::Configuration(format!(
309 "{REDIS_TLS_CA} is required when {REDIS_TLS_POLICY}=custom_ca"
310 )));
311 }
312 if !Path::new(&tls.ca_path).is_absolute() {
313 return Err(TargetError::Configuration(format!("{REDIS_TLS_CA} must be an absolute path")));
314 }
315 }
316 }
317 } else if !tls.ca_path.is_empty() && !Path::new(&tls.ca_path).is_absolute() {
318 return Err(TargetError::Configuration(format!("{REDIS_TLS_CA} must be an absolute path")));
319 }
320
321 Ok(())
322}
323
324pub struct RedisTarget<E>
325where
326 E: PluginEvent,
327{
328 id: TargetID,
329 args: RedisArgs,
330 publisher_client: Arc<parking_lot::Mutex<Client>>,
332 publisher: Arc<Mutex<Option<ConnectionManager>>>,
333 store: Option<Box<dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send + Sync>>,
334 connected: Arc<AtomicBool>,
340 tls_state: Arc<parking_lot::Mutex<super::TargetTlsState>>,
342 tls_adapter: Option<TlsReloadAdapter<Client>>,
346 delivery_counters: Arc<TargetDeliveryCounters>,
347 _phantom: std::marker::PhantomData<E>,
348}
349
350impl<E> RedisTarget<E>
351where
352 E: PluginEvent,
353{
354 #[instrument(skip(args), fields(target_id_as_string = %id))]
355 pub fn new(id: String, args: RedisArgs) -> Result<Self, TargetError> {
356 args.validate()?;
357
358 let target_id = TargetID::new(id, ChannelTargetType::Redis.as_str().to_string());
359 let publisher_client = build_redis_client(&args)?;
360
361 let queue_store = open_target_queue_store(
362 &args.queue_dir,
363 args.queue_limit,
364 args.target_type,
365 ChannelTargetType::Redis.as_str(),
366 &target_id,
367 "Failed to open store for Redis target",
368 )?;
369
370 info!(target_id = %target_id, "Redis target created");
371 Ok(Self {
372 id: target_id,
373 args,
374 publisher_client: Arc::new(parking_lot::Mutex::new(publisher_client)),
375 publisher: Arc::new(Mutex::new(None)),
376 store: queue_store,
377 connected: Arc::new(AtomicBool::new(false)),
378 tls_state: Arc::new(parking_lot::Mutex::new(super::TargetTlsState::default())),
379 tls_adapter: None,
380 delivery_counters: Arc::new(TargetDeliveryCounters::default()),
381 _phantom: std::marker::PhantomData,
382 })
383 }
384
385 pub fn clone_box(&self) -> Box<dyn Target<E> + Send + Sync> {
386 Box::new(Self {
387 id: self.id.clone(),
388 args: self.args.clone(),
389 publisher_client: Arc::clone(&self.publisher_client),
390 publisher: Arc::clone(&self.publisher),
391 store: self.store.as_ref().map(|s| s.boxed_clone()),
392 connected: Arc::clone(&self.connected),
393 tls_state: Arc::clone(&self.tls_state),
394 tls_adapter: self.tls_adapter.clone(),
395 delivery_counters: Arc::clone(&self.delivery_counters),
396 _phantom: std::marker::PhantomData,
397 })
398 }
399
400 async fn get_or_create_publisher(&self) -> Result<ConnectionManager, TargetError> {
401 if let Some(adapter) = &self.tls_adapter {
403 let client: Client = (*adapter.current_material()).clone();
404
405 *self.publisher_client.lock() = client.clone();
407
408 let manager = client
409 .get_connection_manager_lazy(build_redis_connection_manager_config(&self.args))
410 .map_err(map_redis_error)?;
411
412 *self.publisher.lock().await = Some(manager.clone());
413 return Ok(manager);
414 }
415
416 let secure_scheme = matches!(self.args.url.scheme(), "rediss" | "valkeys");
418 if secure_scheme {
419 let next_fingerprint = super::build_target_tls_fingerprint(
420 &self.args.tls.ca_path,
421 &self.args.tls.client_cert_path,
422 &self.args.tls.client_key_path,
423 )
424 .await?;
425 let tls_changed = {
426 let tls_state_guard = self.tls_state.lock();
427 tls_state_guard.needs_update(&next_fingerprint)
428 };
429 if tls_changed {
430 let new_client = build_redis_client(&self.args)?;
431 *self.publisher_client.lock() = new_client;
432 self.invalidate_cached_publisher().await;
433 self.tls_state.lock().refresh(next_fingerprint);
434 }
435 }
436
437 let mut guard = self.publisher.lock().await;
438 if let Some(manager) = guard.clone() {
439 return Ok(manager);
440 }
441
442 let client = self.publisher_client.lock().clone();
443 let manager = client
444 .get_connection_manager_lazy(build_redis_connection_manager_config(&self.args))
445 .map_err(map_redis_error)?;
446
447 *guard = Some(manager.clone());
448 Ok(manager)
449 }
450
451 async fn invalidate_cached_publisher(&self) {
452 *self.publisher.lock().await = None;
456 }
457
458 async fn ensure_publisher_ready(&self) -> Result<(), TargetError> {
459 let mut publisher = self.get_or_create_publisher().await?;
460 match cmd("PING").query_async::<String>(&mut publisher).await {
461 Ok(_) => Ok(()),
462 Err(err) => {
463 let mapped = map_redis_error(err);
464 invalidate_cache_on_connectivity_error(&mapped, || self.invalidate_cached_publisher()).await;
465 Err(mapped)
466 }
467 }
468 }
469
470 async fn init_inner(&self) -> Result<(), TargetError> {
471 if let Err(err) = self.ensure_publisher_ready().await {
472 self.connected.store(false, Ordering::SeqCst);
473 return Err(err);
474 }
475 self.connected.store(true, Ordering::SeqCst);
476 Ok(())
477 }
478
479 #[instrument(skip(self, body, meta), fields(target_id = %self.id))]
480 async fn send_body(&self, body: Vec<u8>, meta: &QueuedPayloadMeta) -> Result<(), TargetError> {
481 debug!(
482 target = %self.id,
483 bucket = %meta.bucket_name,
484 object = %meta.object_name,
485 event = %meta.event_name,
486 payload_len = body.len(),
487 channel = %self.args.channel,
488 "Sending Redis payload"
489 );
490
491 let result = with_delivery_deadline(redis_total_delivery_timeout(&self.args), "Redis delivery", async {
492 let mut attempt = 0usize;
493 let mut last_error = None;
494 while attempt < self.args.max_retry_attempts {
495 attempt += 1;
496
497 let connection_timeout = self.args.connection_timeout.unwrap_or(REDIS_CONNECTION_TIMEOUT_DEFAULT);
498 let mut publisher = match with_delivery_deadline(
499 connection_timeout,
500 "Redis connection",
501 self.get_or_create_publisher(),
502 )
503 .await
504 {
505 Ok(publisher) => publisher,
506 Err(err) => {
507 invalidate_cache_on_connectivity_error(&err, || self.invalidate_cached_publisher()).await;
508 return Err(err);
509 }
510 };
511 let response_timeout = self.args.response_timeout.unwrap_or(REDIS_RESPONSE_TIMEOUT_DEFAULT);
512 match with_delivery_deadline(response_timeout, "Redis publish response", async {
513 publisher
514 .publish::<_, _, i64>(self.args.channel.as_str(), body.as_slice())
515 .await
516 .map_err(map_redis_error)
517 })
518 .await
519 {
520 Ok(receiver_count) => {
521 if receiver_count == 0 {
527 warn!(
528 target_id = %self.id,
529 channel = %self.args.channel,
530 "Redis PUBLISH reached 0 subscribers; the event was not received by any consumer (pub/sub is best-effort)"
531 );
532 }
533 debug!(
534 target_id = %self.id,
535 channel = %self.args.channel,
536 attempt,
537 receiver_count,
538 "Event published to Redis channel"
539 );
540 self.delivery_counters.record_success();
541 return Ok(());
542 }
543 Err(mapped) => {
544 invalidate_cache_on_connectivity_error(&mapped, || self.invalidate_cached_publisher()).await;
545
546 warn!(
547 target_id = %self.id,
548 channel = %self.args.channel,
549 attempt,
550 max_attempts = self.args.max_retry_attempts,
551 error = %mapped,
552 "Redis publish attempt failed"
553 );
554
555 if !is_connectivity_error(&mapped) || attempt >= self.args.max_retry_attempts {
556 last_error = Some(mapped);
557 break;
558 }
559
560 last_error = Some(mapped);
561 tokio::time::sleep(compute_retry_delay(
562 attempt,
563 self.args.min_retry_delay.unwrap_or(Duration::from_millis(100)),
564 self.args.max_retry_delay.unwrap_or(Duration::from_secs(2)),
565 ))
566 .await;
567 }
568 }
569 }
570
571 Err(last_error.unwrap_or(TargetError::Unknown(
572 "Redis publish failed without a captured error".to_string(),
573 )))
574 })
575 .await;
576
577 if let Err(err) = &result {
578 invalidate_cache_on_connectivity_error(err, || self.invalidate_cached_publisher()).await;
579 self.connected.store(false, Ordering::SeqCst);
580 }
581 result
582 }
583}
584
585#[async_trait]
586impl<E> Target<E> for RedisTarget<E>
587where
588 E: PluginEvent,
589{
590 fn id(&self) -> TargetID {
591 self.id.clone()
592 }
593
594 async fn is_active(&self) -> Result<bool, TargetError> {
595 if !self.is_enabled() {
596 return Ok(false);
597 }
598
599 match tokio::time::timeout(REDIS_CONNECTION_TIMEOUT_DEFAULT, self.ensure_publisher_ready()).await {
605 Ok(Ok(())) => {
606 self.connected.store(true, Ordering::SeqCst);
607 Ok(true)
608 }
609 Ok(Err(err)) => {
610 mark_target_disconnected_on_connectivity_error(&self.connected, &err);
611 Err(err)
612 }
613 Err(_) => {
614 let timeout_err = TargetError::Timeout("Redis connection timed out".to_string());
615 invalidate_cache_on_connectivity_error(&timeout_err, || self.invalidate_cached_publisher()).await;
616 mark_target_disconnected_on_connectivity_error(&self.connected, &timeout_err);
617 Err(timeout_err)
618 }
619 }
620 }
621
622 async fn save(&self, event: Arc<EntityTarget<E>>) -> Result<(), TargetError> {
623 let queued = match build_queued_payload(event.as_ref()) {
624 Ok(queued) => queued,
625 Err(err) => {
626 self.delivery_counters.record_final_failure();
627 return Err(err);
628 }
629 };
630
631 if let Some(store) = &self.store {
632 if let Err(e) = persist_queued_payload_to_store(store.as_ref(), &queued) {
633 self.delivery_counters.record_final_failure();
634 return Err(e);
635 }
636
637 debug!(target_id = %self.id, "Event saved to store for Redis target");
638 Ok(())
639 } else {
640 if !self.is_enabled() {
641 return Err(TargetError::Disabled);
642 }
643
644 if let Err(err) = self.init_inner().await {
645 self.delivery_counters.record_final_failure();
646 return Err(err);
647 }
648
649 if let Err(err) = self.send_body(queued.body, &queued.meta).await {
650 self.delivery_counters.record_final_failure();
651 return Err(err);
652 }
653
654 Ok(())
655 }
656 }
657
658 async fn send_raw_from_store(&self, key: Key, body: Vec<u8>, meta: QueuedPayloadMeta) -> Result<(), TargetError> {
659 debug!(target_id = %self.id, ?key, "Attempting to send queued payload from Redis store");
660
661 if !self.is_enabled() {
662 return Err(TargetError::Disabled);
663 }
664
665 if let Err(err) = self.init_inner().await {
666 if is_connectivity_error(&err) {
667 warn!(target_id = %self.id, error = %err, "Redis target not ready; queued event remains in store");
668 }
669 return Err(err);
670 }
671
672 if let Err(err) = self.send_body(body, &meta).await {
673 if is_connectivity_error(&err) {
674 warn!(target_id = %self.id, error = %err, "Failed to send Redis event from store: target not connected. Event remains queued.");
675 }
676 return Err(err);
677 }
678
679 debug!(target_id = %self.id, ?key, "Queued Redis payload sent successfully");
680 Ok(())
681 }
682
683 async fn close(&self) -> Result<(), TargetError> {
684 self.invalidate_cached_publisher().await;
685 self.tls_state.lock().reset();
686 self.connected.store(false, Ordering::SeqCst);
687 info!(target_id = %self.id, "Redis target closed");
688 Ok(())
689 }
690
691 fn store(&self) -> Option<&(dyn Store<QueuedPayload, Error = StoreError, Key = Key> + Send + Sync)> {
692 self.store.as_deref()
693 }
694
695 fn clone_dyn(&self) -> Box<dyn Target<E> + Send + Sync> {
696 self.clone_box()
697 }
698
699 async fn init(&self) -> Result<(), TargetError> {
700 if !self.is_enabled() {
701 return Ok(());
702 }
703 self.init_inner().await
704 }
705
706 fn is_enabled(&self) -> bool {
707 self.args.enable
708 }
709
710 fn delivery_snapshot(&self) -> TargetDeliverySnapshot {
711 self.delivery_counters.snapshot(
712 self.store.as_deref().map_or(0, |store| store.len() as u64),
713 0,
715 )
716 }
717
718 fn record_final_failure(&self) {
719 self.delivery_counters.record_final_failure();
720 }
721}
722
723pub(crate) fn build_redis_client(args: &RedisArgs) -> Result<Client, TargetError> {
724 let mut url = args.url.clone();
725 if args.tls.allow_insecure {
726 url.set_fragment(Some("insecure"));
727 warn!(
731 target = %redact_redis_url(&args.url),
732 "Redis tls_allow_insecure is enabled: server certificate verification is DISABLED (insecure). Use for testing only."
733 );
734 }
735
736 let mut connection_info: ConnectionInfo = url.into_connection_info().map_err(map_redis_error)?;
737
738 let base_redis = connection_info.redis_settings().clone();
739
740 let mut redis_settings = base_redis.clone().set_lib_name("rustfs-targets", env!("CARGO_PKG_VERSION"));
741
742 if let Some(username) = args.username.as_deref().filter(|value| !value.is_empty()) {
743 if base_redis.username().is_some_and(|base| base != username) {
744 warn!(url_username = ?base_redis.username(), arg_username = %username, "Redis target protocol username from URL is being overridden");
745 }
746 redis_settings = redis_settings.set_username(username);
747 }
748 if let Some(password) = args.password.as_deref().filter(|value| !value.is_empty()) {
749 if base_redis.password().is_some() {
750 warn!("RedisArgs.password overrides password from Redis URL");
751 }
752 redis_settings = redis_settings.set_password(password);
753 }
754
755 let mut tcp_settings = TcpSettings::default().set_nodelay(true);
756 #[cfg(not(target_family = "wasm"))]
757 {
758 if !args.keep_alive.is_zero() {
759 tcp_settings = tcp_settings.set_keepalive(socket2::TcpKeepalive::new().with_time(args.keep_alive));
760 }
761 }
762
763 connection_info = connection_info
764 .set_redis_settings(redis_settings)
765 .set_tcp_settings(tcp_settings);
766
767 let secure_scheme = matches!(args.url.scheme(), "rediss" | "valkeys");
768 if secure_scheme {
769 super::ensure_rustls_provider_installed();
770 let tls_certs = TlsCertificates {
771 client_tls: read_client_tls(&args.tls)?,
772 root_cert: read_root_cert(&args.tls)?,
773 };
774 Client::build_with_tls(connection_info, tls_certs).map_err(map_redis_error)
775 } else {
776 Client::open(connection_info).map_err(map_redis_error)
777 }
778}
779
780pub(crate) fn build_redis_connection_manager_config(args: &RedisArgs) -> ConnectionManagerConfig {
781 let mut config = ConnectionManagerConfig::new();
782
783 if let Some(reconnect_retry_attempts) = args.reconnect_retry_attempts {
784 config = config.set_number_of_retries(reconnect_retry_attempts);
785 }
786 if let Some(min_retry_delay) = args.min_retry_delay {
787 config = config.set_min_delay(min_retry_delay);
788 }
789 if let Some(max_retry_delay) = args.max_retry_delay {
790 config = config.set_max_delay(max_retry_delay);
791 }
792 if let Some(connection_timeout) = args.connection_timeout {
793 config = config.set_connection_timeout(Some(connection_timeout));
794 }
795 if let Some(response_timeout) = args.response_timeout {
796 config = config.set_response_timeout(Some(response_timeout));
797 }
798 if let Some(pipeline_buffer_size) = args.pipeline_buffer_size {
799 config = config.set_pipeline_buffer_size(pipeline_buffer_size);
800 }
801
802 config
803}
804
805pub(crate) async fn ping_redis_server(client: &Client, args: &RedisArgs) -> Result<(), TargetError> {
806 let config = build_redis_connection_manager_config(args);
807 let mut conn = client
808 .get_connection_manager_with_config(config)
809 .await
810 .map_err(map_redis_error)?;
811
812 cmd("PING").query_async::<String>(&mut conn).await.map_err(map_redis_error)?;
813
814 Ok(())
815}
816
817fn read_client_tls(tls: &RedisTlsConfig) -> Result<Option<ClientTlsConfig>, TargetError> {
818 if tls.client_cert_path.is_empty() {
819 return Ok(None);
820 }
821
822 let client_cert = std::fs::read(&tls.client_cert_path)
823 .map_err(|e| TargetError::Configuration(format!("Failed to read Redis client cert: {e}")))?;
824 let client_key = std::fs::read(&tls.client_key_path)
825 .map_err(|e| TargetError::Configuration(format!("Failed to read Redis client key: {e}")))?;
826
827 Ok(Some(ClientTlsConfig { client_cert, client_key }))
828}
829
830fn read_root_cert(tls: &RedisTlsConfig) -> Result<Option<Vec<u8>>, TargetError> {
831 if tls.ca_path.is_empty() {
832 return Ok(None);
833 }
834
835 let pem =
836 std::fs::read(&tls.ca_path).map_err(|e| TargetError::Configuration(format!("Failed to read Redis root CA cert: {e}")))?;
837 let mut reader = BufReader::new(pem.as_slice());
838 let certs_der = CertificateDer::pem_reader_iter(&mut reader)
839 .collect::<Result<Vec<_>, _>>()
840 .map_err(|e| TargetError::Configuration(format!("Failed to parse Redis root CA cert: {e}")))?;
841
842 if certs_der.is_empty() {
843 return Err(TargetError::Configuration(
844 "Redis root CA cert did not contain any parsable certificates".to_string(),
845 ));
846 }
847
848 Ok(Some(pem))
849}
850
851fn map_redis_error(err: RedisError) -> TargetError {
852 use redis::ErrorKind;
853
854 match err.kind() {
855 ErrorKind::AuthenticationFailed => TargetError::Authentication(err.to_string()),
856 ErrorKind::RESP3NotSupported => TargetError::Initialization(err.to_string()),
857 ErrorKind::InvalidClientConfig => TargetError::Configuration(err.to_string()),
858 ErrorKind::Io if err.is_timeout() => TargetError::Timeout(err.to_string()),
859 ErrorKind::Io if err.is_connection_dropped() || err.is_connection_refusal() => TargetError::NotConnected,
860 ErrorKind::Io => TargetError::Network(err.to_string()),
861 _ if err.is_unrecoverable_error() => TargetError::NotConnected,
862 _ => TargetError::Request(err.to_string()),
863 }
864}
865
866fn compute_retry_delay(attempt: usize, min_delay: Duration, max_delay: Duration) -> Duration {
867 let shift = attempt.saturating_sub(1).min(16) as u32;
868 let factor = 1u32 << shift;
869 min_delay.saturating_mul(factor).min(max_delay)
870}
871
872#[async_trait]
877impl<E> ReloadableTargetTls for RedisTarget<E>
878where
879 E: PluginEvent,
880{
881 type Material = Client;
882
883 fn tls_input_set(&self) -> TargetTlsInputSet {
884 TargetTlsInputSet {
885 ca_path: self.args.tls.ca_path.clone(),
886 client_cert_path: self.args.tls.client_cert_path.clone(),
887 client_key_path: self.args.tls.client_key_path.clone(),
888 target_label: format!("redis:{}", self.id.id),
889 }
890 }
891
892 async fn build_tls_material(&self) -> Result<Self::Material, TargetError> {
893 build_redis_client(&self.args)
894 }
895
896 async fn apply_tls_material(
897 &self,
898 _generation: TargetTlsGeneration,
899 material: Arc<Self::Material>,
900 _mode: ReloadApplyMode,
901 ) -> Result<(), TargetError> {
902 *self.publisher_client.lock() = (*material).clone();
903 self.invalidate_cached_publisher().await;
904 Ok(())
905 }
906
907 async fn validate_tls_files(&self) -> Result<(), TargetError> {
908 validate_tls_material(&self.args.tls.ca_path, &self.args.tls.client_cert_path, &self.args.tls.client_key_path)
909 }
910}
911
912#[cfg(test)]
913mod tests {
914 use super::*;
915 use redis::ProtocolVersion;
916 use std::sync::atomic::Ordering;
917 use tokio::io::{AsyncReadExt, AsyncWriteExt};
918 use tokio::net::TcpListener;
919
920 fn absolute_test_path(path: &str) -> String {
921 std::env::temp_dir().join(path).to_string_lossy().into_owned()
922 }
923
924 fn base_args() -> RedisArgs {
925 RedisArgs {
926 enable: true,
927 url: Url::parse("redis://127.0.0.1:6379").unwrap(),
928 channel: "rustfs-events".to_string(),
929 username: None,
930 password: None,
931 tls: RedisTlsConfig::default(),
932 keep_alive: Duration::from_secs(15),
933 queue_dir: String::new(),
934 queue_limit: 0,
935 max_retry_attempts: 3,
936 reconnect_retry_attempts: None,
937 min_retry_delay: None,
938 max_retry_delay: None,
939 connection_timeout: None,
940 response_timeout: None,
941 pipeline_buffer_size: None,
942 target_type: TargetType::NotifyEvent,
943 }
944 }
945
946 #[test]
947 fn validate_rejects_empty_channel() {
948 let args = RedisArgs {
949 channel: String::new(),
950 ..base_args()
951 };
952 assert!(args.validate().is_err());
953 }
954
955 #[test]
956 fn validate_accepts_embedded_credentials_in_url() {
957 let url = Url::parse("redis://user:pass@127.0.0.1:6379").unwrap();
958 assert!(validate_redis_url(&url).is_ok());
959 }
960
961 #[test]
962 fn validate_rejects_relative_queue_dir() {
963 let args = RedisArgs {
964 queue_dir: "relative/path".to_string(),
965 ..base_args()
966 };
967 assert!(args.validate().is_err());
968 }
969
970 #[test]
971 fn validate_rejects_zero_connection_timeout() {
972 let args = RedisArgs {
973 connection_timeout: Some(Duration::ZERO),
974 ..base_args()
975 };
976
977 assert!(args.validate().is_err());
978 }
979
980 #[test]
981 fn validate_rejects_zero_response_timeout() {
982 let args = RedisArgs {
983 response_timeout: Some(Duration::ZERO),
984 ..base_args()
985 };
986
987 assert!(args.validate().is_err());
988 }
989
990 #[test]
991 fn validate_accepts_custom_ca_tls_policy() {
992 let args = RedisArgs {
993 url: Url::parse("rediss://127.0.0.1:6379").unwrap(),
994 tls: RedisTlsConfig {
995 policy: Some(RedisTlsPolicy::CustomCa),
996 ca_path: absolute_test_path("redis-ca.pem"),
997 ..RedisTlsConfig::default()
998 },
999 ..base_args()
1000 };
1001 assert!(args.validate().is_ok());
1002 }
1003
1004 #[test]
1005 fn debug_redacts_passwords_from_url_and_args() {
1006 let args = RedisArgs {
1007 url: Url::parse("redis://user:secret@127.0.0.1:6379/0").unwrap(),
1008 password: Some("override-secret".to_string()),
1009 ..base_args()
1010 };
1011
1012 let rendered = format!("{args:?}");
1013 assert!(!rendered.contains("secret"), "url password leaked: {rendered}");
1014 assert!(!rendered.contains("override-secret"), "args password leaked: {rendered}");
1015 assert!(rendered.contains("redis://user:***@127.0.0.1:6379/0"));
1016 assert!(rendered.contains("\"***REDACTED***\""));
1017 }
1018
1019 #[test]
1020 fn validate_rejects_insecure_tls_for_non_secure_scheme() {
1021 let args = RedisArgs {
1022 tls: RedisTlsConfig {
1023 allow_insecure: true,
1024 ..RedisTlsConfig::default()
1025 },
1026 ..base_args()
1027 };
1028 assert!(args.validate().is_err());
1029 }
1030
1031 #[test]
1032 fn build_redis_client_preserves_url_auth_when_args_are_none() {
1033 let args = RedisArgs {
1034 url: Url::parse("redis://user:pass@127.0.0.1:6379/2").unwrap(),
1035 ..base_args()
1036 };
1037
1038 let client = build_redis_client(&args).expect("client should build");
1039 let info = client.get_connection_info();
1040 let redis = info.redis_settings();
1041
1042 assert_eq!(redis.username(), Some("user"));
1043 assert_eq!(redis.password(), Some("pass"));
1044 assert_eq!(redis.db(), 2);
1045 }
1046
1047 #[test]
1048 fn build_redis_client_overrides_url_auth_when_args_are_set() {
1049 let args = RedisArgs {
1050 url: Url::parse("redis://user:pass@127.0.0.1:6379/2").unwrap(),
1051 username: Some("override-user".to_string()),
1052 password: Some("override-pass".to_string()),
1053 ..base_args()
1054 };
1055
1056 let client = build_redis_client(&args).expect("client should build");
1057 let redis = client.get_connection_info().redis_settings();
1058
1059 assert_eq!(redis.username(), Some("override-user"));
1060 assert_eq!(redis.password(), Some("override-pass"));
1061 assert_eq!(redis.db(), 2);
1062 }
1063
1064 #[test]
1065 fn build_redis_client_preserves_url_protocol_when_args_do_not_override_it() {
1066 let args = RedisArgs {
1067 url: Url::parse("redis://127.0.0.1:6379/?protocol=resp3").unwrap(),
1068 ..base_args()
1069 };
1070
1071 let client = build_redis_client(&args).expect("client should build");
1072 let redis = client.get_connection_info().redis_settings();
1073
1074 assert_eq!(redis.protocol(), ProtocolVersion::RESP3);
1075 }
1076
1077 #[test]
1078 fn build_redis_client_enables_insecure_tls_when_requested() {
1079 let args = RedisArgs {
1080 url: Url::parse("rediss://127.0.0.1:6379").unwrap(),
1081 tls: RedisTlsConfig {
1082 allow_insecure: true,
1083 ..RedisTlsConfig::default()
1084 },
1085 ..base_args()
1086 };
1087
1088 let client = build_redis_client(&args).expect("client should build");
1089 match client.get_connection_info().addr() {
1090 redis::ConnectionAddr::TcpTls { insecure, .. } => assert!(*insecure),
1091 other => panic!("expected TLS address, got {other:?}"),
1092 }
1093 }
1094
1095 #[tokio::test]
1096 async fn invalidate_cached_publisher_keeps_connected_state() {
1097 let target = RedisTarget::<String>::new("redis:test".to_string(), base_args()).expect("target should build");
1098 target.connected.store(true, Ordering::SeqCst);
1099
1100 target.invalidate_cached_publisher().await;
1101
1102 assert!(target.connected.load(Ordering::SeqCst));
1103 assert!(target.publisher.lock().await.is_none());
1104 }
1105
1106 #[tokio::test]
1107 async fn is_active_succeeds_when_ping_returns_pong() {
1108 let listener = match TcpListener::bind("127.0.0.1:0").await {
1109 Ok(listener) => listener,
1110 Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
1111 Err(err) => panic!("bind fake redis: {err}"),
1112 };
1113 let addr = listener.local_addr().expect("listener addr");
1114 tokio::spawn(run_fake_redis_server(listener, false));
1115
1116 let mut args = base_args();
1117 args.url = Url::parse(&format!("redis://{}:{}/0", addr.ip(), addr.port())).unwrap();
1118
1119 let target = RedisTarget::<String>::new("redis:test".to_string(), args).expect("target should build");
1120 target.connected.store(false, Ordering::SeqCst);
1121
1122 assert!(target.is_active().await.expect("ping should succeed"));
1123 assert!(target.connected.load(Ordering::SeqCst));
1124 }
1125
1126 #[tokio::test]
1127 async fn is_active_returns_error_when_ping_fails() {
1128 let listener = match TcpListener::bind("127.0.0.1:0").await {
1129 Ok(listener) => listener,
1130 Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
1131 Err(err) => panic!("bind fake redis: {err}"),
1132 };
1133 let addr = listener.local_addr().expect("listener addr");
1134 tokio::spawn(async move {
1135 loop {
1136 let Ok((socket, _)) = listener.accept().await else {
1137 return;
1138 };
1139 drop(socket);
1140 }
1141 });
1142
1143 let mut args = base_args();
1144 args.url = Url::parse(&format!("redis://{}:{}/0", addr.ip(), addr.port())).unwrap();
1145
1146 let target = RedisTarget::<String>::new("redis:test".to_string(), args).expect("target should build");
1147 target.connected.store(true, Ordering::SeqCst);
1148
1149 let err = target.is_active().await.expect_err("ping should fail");
1150 assert!(matches!(
1151 err,
1152 TargetError::NotConnected | TargetError::Network(_) | TargetError::Timeout(_)
1153 ));
1154 assert!(!target.connected.load(Ordering::SeqCst));
1155 assert!(target.publisher.lock().await.is_none());
1156 }
1157
1158 #[tokio::test]
1159 async fn is_active_returns_false_when_disabled() {
1160 let target = RedisTarget::<String>::new(
1161 "redis:test".to_string(),
1162 RedisArgs {
1163 enable: false,
1164 ..base_args()
1165 },
1166 )
1167 .expect("target should build");
1168
1169 assert!(!target.is_active().await.expect("disabled target should not probe"));
1170 }
1171
1172 #[test]
1173 fn compute_retry_delay_is_bounded() {
1174 let min = Duration::from_millis(100);
1175 let max = Duration::from_secs(2);
1176
1177 assert_eq!(compute_retry_delay(1, min, max), min);
1178 assert!(compute_retry_delay(5, min, max) <= max);
1179 assert_eq!(compute_retry_delay(50, min, max), max);
1180 }
1181
1182 #[test]
1183 fn queued_payload_uses_event_data_in_records() {
1184 let payload = build_queued_payload(&EntityTarget {
1185 object_name: "greeting+file+%282%29.csv".to_string(),
1186 bucket_name: "bucket".to_string(),
1187 event_name: rustfs_s3_types::EventName::ObjectCreatedPut,
1188 data: "payload-data".to_string(),
1189 })
1190 .expect("payload should build");
1191
1192 let value: serde_json::Value = serde_json::from_slice(&payload.body).expect("payload JSON");
1193 assert_eq!(value["Key"], "bucket/greeting file (2).csv");
1194 assert_eq!(value["Records"][0], "payload-data");
1195 }
1196
1197 fn parse_resp_array(input: &[u8]) -> Option<(Vec<String>, usize)> {
1198 if input.first()? != &b'*' {
1199 return None;
1200 }
1201
1202 let mut index = 1;
1203 let len_end = input[index..].windows(2).position(|w| w == b"\r\n")? + index;
1204 let items: usize = std::str::from_utf8(&input[index..len_end]).ok()?.parse().ok()?;
1205 index = len_end + 2;
1206
1207 let mut out = Vec::with_capacity(items);
1208 for _ in 0..items {
1209 if input.get(index)? != &b'$' {
1210 return None;
1211 }
1212 index += 1;
1213 let bulk_end = input[index..].windows(2).position(|w| w == b"\r\n")? + index;
1214 let bulk_len: usize = std::str::from_utf8(&input[index..bulk_end]).ok()?.parse().ok()?;
1215 index = bulk_end + 2;
1216
1217 let data_end = index.checked_add(bulk_len)?;
1218 let data = std::str::from_utf8(input.get(index..data_end)?).ok()?.to_string();
1219 out.push(data);
1220 index = data_end + 2;
1221 }
1222
1223 Some((out, index))
1224 }
1225
1226 async fn run_fake_redis_server(listener: TcpListener, close_first_connection: bool) {
1227 let mut first = close_first_connection;
1228 loop {
1229 let Ok((mut socket, _)) = listener.accept().await else {
1230 return;
1231 };
1232
1233 if first {
1234 first = false;
1235 drop(socket);
1236 continue;
1237 }
1238
1239 tokio::spawn(async move {
1240 let mut buf = vec![0_u8; 4096];
1241 let mut pending = Vec::new();
1242
1243 loop {
1244 let Ok(read) = socket.read(&mut buf).await else {
1245 return;
1246 };
1247 if read == 0 {
1248 return;
1249 }
1250
1251 pending.extend_from_slice(&buf[..read]);
1252
1253 while let Some((command, consumed)) = parse_resp_array(&pending) {
1254 pending.drain(..consumed);
1255 let response = match command.first().map(|s| s.as_str()) {
1256 Some("PING") => b"+PONG\r\n".as_slice(),
1257 Some("PUBLISH") => b":1\r\n".as_slice(),
1258 Some("CLIENT") => b"+OK\r\n".as_slice(),
1259 Some("AUTH") => b"+OK\r\n".as_slice(),
1260 Some("SELECT") => b"+OK\r\n".as_slice(),
1261 Some("HELLO") => b"%1\r\n+server\r\n+redis\r\n".as_slice(),
1262 _ => b"+OK\r\n".as_slice(),
1263 };
1264
1265 if socket.write_all(response).await.is_err() {
1266 return;
1267 }
1268 }
1269 }
1270 });
1271 }
1272 }
1273
1274 #[tokio::test]
1275 async fn send_body_keeps_connected_true_when_retryable_error_eventually_recovers() {
1276 let listener = match TcpListener::bind("127.0.0.1:0").await {
1277 Ok(listener) => listener,
1278 Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
1279 Err(err) => panic!("bind fake redis: {err}"),
1280 };
1281 let addr = listener.local_addr().expect("listener addr");
1282 tokio::spawn(run_fake_redis_server(listener, true));
1283
1284 let mut args = base_args();
1285 args.url = Url::parse(&format!("redis://{}:{}/0", addr.ip(), addr.port())).unwrap();
1286 args.max_retry_attempts = 3;
1287 args.reconnect_retry_attempts = Some(0);
1288 args.min_retry_delay = Some(Duration::from_millis(10));
1289 args.max_retry_delay = Some(Duration::from_millis(20));
1290
1291 let target = RedisTarget::<String>::new("redis:test".to_string(), args).expect("target should build");
1292 let meta = QueuedPayloadMeta::new(
1293 rustfs_s3_types::EventName::ObjectCreatedPut,
1294 "bucket".to_string(),
1295 "object".to_string(),
1296 "application/json",
1297 2,
1298 );
1299
1300 target.connected.store(true, Ordering::SeqCst);
1301 target
1302 .send_body(b"{}".to_vec(), &meta)
1303 .await
1304 .expect("eventual retry should succeed");
1305
1306 assert!(target.connected.load(Ordering::SeqCst));
1307 assert_eq!(target.delivery_snapshot().total_messages, 1);
1308 }
1309
1310 #[tokio::test(start_paused = true)]
1311 async fn delivery_budget_respects_response_timeout_longer_than_sixty_seconds() {
1312 let mut args = base_args();
1313 args.max_retry_attempts = 1;
1314 args.response_timeout = Some(Duration::from_secs(90));
1315 let manager_config = build_redis_connection_manager_config(&args);
1316
1317 assert_eq!(manager_config.response_timeout(), Some(Duration::from_secs(90)));
1318 assert_eq!(redis_total_delivery_timeout(&args), Duration::from_secs(97));
1319 with_delivery_deadline(redis_total_delivery_timeout(&args), "Redis delivery", async {
1320 tokio::time::sleep(Duration::from_secs(70)).await;
1321 Ok::<_, TargetError>(())
1322 })
1323 .await
1324 .expect("the configured delivery budget must not impose a fixed sixty-second cap");
1325 }
1326
1327 #[tokio::test]
1328 async fn send_body_sets_connected_false_after_retry_exhaustion() {
1329 let listener = match TcpListener::bind("127.0.0.1:0").await {
1330 Ok(listener) => listener,
1331 Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
1332 Err(err) => panic!("bind fake redis: {err}"),
1333 };
1334 let addr = listener.local_addr().expect("listener addr");
1335 tokio::spawn(async move {
1336 loop {
1337 let Ok((socket, _)) = listener.accept().await else {
1338 return;
1339 };
1340 drop(socket);
1341 }
1342 });
1343
1344 let mut args = base_args();
1345 args.url = Url::parse(&format!("redis://{}:{}/0", addr.ip(), addr.port())).unwrap();
1346 args.max_retry_attempts = 2;
1347 args.reconnect_retry_attempts = Some(0);
1348 args.min_retry_delay = Some(Duration::from_millis(10));
1349 args.max_retry_delay = Some(Duration::from_millis(20));
1350
1351 let target = RedisTarget::<String>::new("redis:test".to_string(), args).expect("target should build");
1352 let meta = QueuedPayloadMeta::new(
1353 rustfs_s3_types::EventName::ObjectCreatedPut,
1354 "bucket".to_string(),
1355 "object".to_string(),
1356 "application/json",
1357 2,
1358 );
1359
1360 let err = target
1361 .send_body(b"{}".to_vec(), &meta)
1362 .await
1363 .expect_err("all retries should fail");
1364 assert!(matches!(
1365 err,
1366 TargetError::NotConnected | TargetError::Network(_) | TargetError::Timeout(_)
1367 ));
1368 assert!(!target.connected.load(Ordering::SeqCst));
1369 assert_eq!(target.delivery_snapshot().total_messages, 0);
1370 }
1371
1372 #[tokio::test]
1373 async fn send_raw_from_store_failure_does_not_count_as_success() {
1374 let listener = match TcpListener::bind("127.0.0.1:0").await {
1375 Ok(listener) => listener,
1376 Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
1377 Err(err) => panic!("bind fake redis: {err}"),
1378 };
1379 let addr = listener.local_addr().expect("listener addr");
1380 tokio::spawn(async move {
1381 loop {
1382 let Ok((socket, _)) = listener.accept().await else {
1383 return;
1384 };
1385 drop(socket);
1386 }
1387 });
1388
1389 let mut args = base_args();
1390 args.url = Url::parse(&format!("redis://{}:{}/0", addr.ip(), addr.port())).unwrap();
1391 args.max_retry_attempts = 1;
1392 args.reconnect_retry_attempts = Some(0);
1393
1394 let target = RedisTarget::<String>::new("redis:test".to_string(), args).expect("target should build");
1395 let meta = QueuedPayloadMeta::new(
1396 rustfs_s3_types::EventName::ObjectCreatedPut,
1397 "bucket".to_string(),
1398 "object".to_string(),
1399 "application/json",
1400 2,
1401 );
1402
1403 let err = target
1404 .send_raw_from_store(
1405 Key {
1406 name: "key".to_string(),
1407 extension: String::new(),
1408 item_count: 1,
1409 compress: false,
1410 },
1411 b"{}".to_vec(),
1412 meta,
1413 )
1414 .await
1415 .expect_err("send from store should fail");
1416
1417 assert!(matches!(
1418 err,
1419 TargetError::NotConnected | TargetError::Network(_) | TargetError::Timeout(_)
1420 ));
1421 assert_eq!(target.delivery_snapshot().total_messages, 0);
1422 }
1423}