Skip to main content

rustfs_targets/target/
redis.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::{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    /// Whether the target is enabled
130    pub enable: bool,
131    /// The Redis server URL in format: `{redis|rediss|valkey|valkeys}://[<username>][:<password>@]<hostname>[:port][/<db>]`
132    pub url: Url,
133    /// The Redis pub/sub channel to publish to
134    pub channel: String,
135    /// The username for the Redis connection (leave it empty if you parse with url)
136    pub username: Option<String>,
137    /// The password for the Redis connection (leave it empty if you parse with url)
138    pub password: Option<String>,
139    /// TLS configuration
140    pub tls: RedisTlsConfig,
141    /// The keep alive interval
142    pub keep_alive: Duration,
143    /// The directory to store events in case of failure
144    pub queue_dir: String,
145    /// The maximum number of events to store
146    pub queue_limit: u64,
147    /// Maximum number of synchronous publish retries per payload
148    pub max_retry_attempts: usize,
149    /// Maximum number of reconnect retries in the underlying connection manager (6 if not provided)
150    pub reconnect_retry_attempts: Option<usize>,
151    /// Minimum retry delay between publish retry attempts (100ms if not provided)
152    pub min_retry_delay: Option<Duration>,
153    /// Maximum retry delay between publish retry attempts (2s if not provided)
154    pub max_retry_delay: Option<Duration>,
155    /// Timeout for establishing a Redis connection (5s if not provided)
156    pub connection_timeout: Option<Duration>,
157    /// Timeout for command responses (5s if not provided)
158    pub response_timeout: Option<Duration>,
159    /// Internal command buffer size for the multiplexed connection (50 if not provided)
160    pub pipeline_buffer_size: Option<usize>,
161    /// the target type
162    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    /// Redis client, wrapped in a lock so TLS hot-reload can atomically replace it.
331    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    /// Business-level liveness flag.
335    ///
336    /// We only flip this to `false` on final/terminal failure paths (for example: init failed,
337    /// publish exhausted retries, or the target was explicitly closed). Temporary reconnectable
338    /// errors only invalidate the cached publisher so that a later request can lazily rebuild it.
339    connected: Arc<AtomicBool>,
340    /// TLS fingerprint tracking for hot reload (inline fallback path).
341    tls_state: Arc<parking_lot::Mutex<super::TargetTlsState>>,
342    /// Adapter that bridges this target to the TLS reload coordinator.
343    /// When `Some`, the target uses coordinator-managed material; when `None`,
344    /// it falls back to inline fingerprint-based change detection.
345    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        // Adapter-managed path: use the material directly from the TLS reload adapter.
402        if let Some(adapter) = &self.tls_adapter {
403            let client: Client = (*adapter.current_material()).clone();
404
405            // Ensure the client is also stored locally so close() can drain it.
406            *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        // Inline fingerprint fallback path (no coordinator).
417        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        // Intentionally does not touch `connected`: invalidating the current manager only means
453        // "recreate the publisher on the next attempt", not "this target is now definitively
454        // inactive". That distinction preserves the business semantics of `is_active()`.
455        *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                        // PUBLISH returns the number of subscribers that received the
522                        // message. Redis pub/sub is best-effort: with zero subscribers
523                        // the event is delivered to no one, yet the durable copy is
524                        // deleted. Warn so operators relying on reliable delivery are
525                        // not silently losing events (backlog#982).
526                        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        // Reuse the cached ConnectionManager for the health probe (via
600        // ensure_publisher_ready) instead of building a brand-new manager — and
601        // thus a fresh TCP+TLS handshake — on every health check (backlog#982).
602        // ensure_publisher_ready already invalidates the cached manager on a
603        // connectivity error so the next attempt rebuilds it.
604        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            // Redis targets record no terminal failures and keep no failed store.
714            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        // Mirror the webhook `skip_tls_verify` warning: this disables Redis
728        // server certificate verification and must only be used for testing
729        // (backlog#982).
730        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/// Coordinated TLS hot-reload implementation for Redis targets.
873///
874/// The coordinator calls these methods on a background poll loop to detect
875/// TLS file changes and rebuild the Redis client without restarting.
876#[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}