// Copyright 2024 RustFS Team // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. use crate::plugin::PluginEvent; use crate::{ StoreError, Target, arn::TargetID, error::TargetError, runtime::tls::{ ReloadableTargetTls, TargetTlsGeneration, TargetTlsInputSet, TlsReloadAdapter, config::ReloadApplyMode, validate_tls_material, }, store::{Key, Store}, target::{ ChannelTargetType, EntityTarget, QueuedPayload, QueuedPayloadMeta, TargetDeliveryCounters, TargetDeliverySnapshot, TargetTlsState, TargetType, build_queued_payload_with_records, build_target_tls_fingerprint, open_target_queue_store, persist_queued_payload_to_store, redacted_secret, }, }; use async_trait::async_trait; use pulsar::{Authentication, Producer, Pulsar, TokioExecutor}; use rustfs_tls_runtime::load_cert_bundle_der_bytes; use std::fmt; use std::path::Path; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use tokio::sync::Mutex as AsyncMutex; use tracing::{info, instrument}; use url::Url; #[derive(Clone)] pub struct PulsarArgs { pub enable: bool, pub broker: String, pub topic: String, pub auth_token: String, pub username: String, pub password: String, pub tls_ca: String, pub tls_allow_insecure: bool, pub tls_hostname_verification: bool, pub queue_dir: String, pub queue_limit: u64, pub target_type: TargetType, } impl fmt::Debug for PulsarArgs { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.debug_struct("PulsarArgs") .field("enable", &self.enable) .field("broker", &self.broker) .field("topic", &self.topic) .field("auth_token", &redacted_secret(&self.auth_token)) .field("username", &self.username) .field("password", &redacted_secret(&self.password)) .field("tls_ca", &self.tls_ca) .field("tls_allow_insecure", &self.tls_allow_insecure) .field("tls_hostname_verification", &self.tls_hostname_verification) .field("queue_dir", &self.queue_dir) .field("queue_limit", &self.queue_limit) .field("target_type", &self.target_type) .finish() } } impl PulsarArgs { pub fn validate(&self) -> Result<(), TargetError> { if !self.enable { return Ok(()); } validate_pulsar_broker(&self.broker)?; if self.topic.trim().is_empty() { return Err(TargetError::Configuration("Pulsar topic cannot be empty".to_string())); } if !self.auth_token.is_empty() && (!self.username.is_empty() || !self.password.is_empty()) { return Err(TargetError::Configuration( "Pulsar supports either auth_token or username/password auth, not both".to_string(), )); } if self.username.is_empty() != self.password.is_empty() { return Err(TargetError::Configuration( "Pulsar username and password must be specified together".to_string(), )); } if !self.tls_ca.is_empty() && !Path::new(&self.tls_ca).is_absolute() { return Err(TargetError::Configuration("Pulsar tls_ca must be an absolute path".to_string())); } if !self.queue_dir.is_empty() && !Path::new(&self.queue_dir).is_absolute() { return Err(TargetError::Configuration("Pulsar queue directory must be an absolute path".to_string())); } let parsed = Url::parse(&self.broker) .map_err(|e| TargetError::Configuration(format!("Invalid Pulsar broker URL: {e} (value: '{}')", self.broker)))?; let tls_enabled = parsed.scheme() == "pulsar+ssl"; if !tls_enabled && (!self.tls_ca.is_empty() || self.tls_allow_insecure || !self.tls_hostname_verification) { return Err(TargetError::Configuration( "Pulsar TLS settings are only allowed with pulsar+ssl brokers".to_string(), )); } Ok(()) } } pub fn validate_pulsar_broker(broker: &str) -> Result { let url = Url::parse(broker) .map_err(|e| TargetError::Configuration(format!("Invalid Pulsar broker URL: {e} (value: '{broker}')")))?; match url.scheme() { "pulsar" | "pulsar+ssl" => {} _ => { return Err(TargetError::Configuration( "Pulsar broker must use pulsar:// or pulsar+ssl://".to_string(), )); } } if !url.username().is_empty() || url.password().is_some() { return Err(TargetError::Configuration( "Pulsar broker URL must not embed username or password".to_string(), )); } if url.host_str().is_none() { return Err(TargetError::Configuration("Pulsar broker is missing host".to_string())); } Ok(url) } pub async fn connect_pulsar(args: &PulsarArgs) -> Result, TargetError> { args.validate()?; let mut builder = Pulsar::builder(args.broker.clone(), TokioExecutor); if !args.auth_token.is_empty() { builder = builder.with_auth(Authentication { name: "token".to_string(), data: args.auth_token.clone().into_bytes(), }); } else if !args.username.is_empty() { builder = builder.with_auth_provider(pulsar::authentication::basic::BasicAuthentication::new(&args.username, &args.password)); } if !args.tls_ca.is_empty() { let certs = load_cert_bundle_der_bytes(&args.tls_ca) .map_err(|e| TargetError::Configuration(format!("Failed to parse Pulsar tls_ca: {e}")))?; if certs.is_empty() { return Err(TargetError::Configuration( "Pulsar tls_ca did not contain any parsable certificates".to_string(), )); } builder = builder .with_certificate_chain_file(&args.tls_ca) .map_err(|e| TargetError::Configuration(format!("Failed to load Pulsar tls_ca: {e}")))?; } builder = builder .with_allow_insecure_connection(args.tls_allow_insecure) .with_tls_hostname_verification_enabled(args.tls_hostname_verification); builder .build() .await .map_err(|e| TargetError::Network(format!("Failed to connect to Pulsar broker: {e}"))) } pub struct PulsarTarget where E: PluginEvent, { id: TargetID, args: PulsarArgs, client: Mutex>>, tls_state: Mutex, /// When set, the coordinator drives TLS reload; inline fingerprint check is skipped. tls_adapter: Option>>, producer: AsyncMutex>>, store: Option + Send + Sync>>, connected: AtomicBool, delivery_counters: Arc, _phantom: std::marker::PhantomData, } impl PulsarTarget where E: PluginEvent, { pub fn clone_box(&self) -> Box + Send + Sync> { Box::new(PulsarTarget:: { id: self.id.clone(), args: self.args.clone(), client: Mutex::new(self.client.lock().unwrap().clone()), tls_state: Mutex::new(self.tls_state.lock().unwrap().clone()), tls_adapter: self.tls_adapter.clone(), producer: AsyncMutex::new(None), store: self.store.as_ref().map(|s| s.boxed_clone()), connected: AtomicBool::new(self.connected.load(Ordering::SeqCst)), delivery_counters: Arc::clone(&self.delivery_counters), _phantom: std::marker::PhantomData, }) } #[instrument(skip(args), fields(target_id_as_string = %id))] pub fn new(id: String, args: PulsarArgs) -> Result { args.validate()?; let target_id = TargetID::new(id, ChannelTargetType::Pulsar.as_str().to_string()); let queue_store = open_target_queue_store( &args.queue_dir, args.queue_limit, args.target_type, ChannelTargetType::Pulsar.as_str(), &target_id, "Failed to open store for Pulsar target", )?; Ok(Self { id: target_id, args, client: Mutex::new(None), tls_state: Mutex::new(TargetTlsState::default()), tls_adapter: None, producer: AsyncMutex::new(None), store: queue_store, connected: AtomicBool::new(false), delivery_counters: Arc::new(TargetDeliveryCounters::default()), _phantom: std::marker::PhantomData, }) } fn clear_cached_client_connection(&self) { self.client.lock().unwrap().take(); } fn clear_cached_client(&self) { self.clear_cached_client_connection(); self.tls_state.lock().unwrap().reset(); } async fn get_or_connect_client(&self) -> Result, TargetError> { // When a TLS reload adapter is attached, it drives client rebuilds // in the background. The inline per-send fingerprint check is skipped. if let Some(adapter) = &self.tls_adapter { let material = adapter.current_material(); { let mut guard = self.client.lock().unwrap(); *guard = Some((*material).clone()); } } else { let next_fingerprint = build_target_tls_fingerprint(&self.args.tls_ca, "", "").await?; let tls_changed = { let tls_state_guard = self.tls_state.lock().unwrap(); tls_state_guard.needs_update(&next_fingerprint) }; if tls_changed { self.clear_cached_client_connection(); self.tls_state.lock().unwrap().refresh(next_fingerprint); } } if let Some(client) = self.client.lock().unwrap().clone() { return Ok(client); } let client = connect_pulsar(&self.args).await?; self.connected.store(true, Ordering::SeqCst); let mut guard = self.client.lock().unwrap(); let shared = guard.get_or_insert_with(|| client.clone()).clone(); Ok(shared) } async fn init_producer(&self) -> Result<(), TargetError> { if self.producer.lock().await.is_some() { return Ok(()); } let client = self.get_or_connect_client().await?; let producer = client .producer() .with_topic(self.args.topic.clone()) .with_name(self.id.id.clone()) .build() .await .map_err(|e| TargetError::Network(format!("Failed to create Pulsar producer: {e}")))?; let mut guard = self.producer.lock().await; if guard.is_none() { *guard = Some(producer); } Ok(()) } fn build_queued_payload(&self, event: &EntityTarget) -> Result { build_queued_payload_with_records(event, vec![event.clone()]) } async fn send_body(&self, body: Vec) -> Result<(), TargetError> { self.init_producer().await?; let mut guard = self.producer.lock().await; let producer = guard .as_mut() .ok_or_else(|| TargetError::Configuration("Pulsar producer not initialized".to_string()))?; let receipt = producer .send_non_blocking(body) .await .map_err(|e| TargetError::Request(format!("Failed to send Pulsar message: {e}")))?; receipt .await .map_err(|e| TargetError::Request(format!("Failed to receive Pulsar receipt: {e}")))?; self.delivery_counters.record_success(); Ok(()) } } /// Coordinated TLS hot-reload implementation for Pulsar targets. /// /// Pulsar only uses a CA certificate (no client cert/key). /// The coordinator calls these methods on a background poll loop to detect /// TLS file changes and rebuild the client without restarting. #[async_trait] impl ReloadableTargetTls for PulsarTarget where E: PluginEvent, { type Material = Pulsar; fn tls_input_set(&self) -> TargetTlsInputSet { TargetTlsInputSet { ca_path: self.args.tls_ca.clone(), client_cert_path: String::new(), client_key_path: String::new(), target_label: format!("pulsar:{}", self.id.id), } } async fn build_tls_material(&self) -> Result { connect_pulsar(&self.args).await } async fn apply_tls_material( &self, _generation: TargetTlsGeneration, material: Arc, _mode: ReloadApplyMode, ) -> Result<(), TargetError> { // Pulsar client is Clone, so we clone from the Arc and store it. { let mut guard = self.client.lock().unwrap(); *guard = Some((*material).clone()); } // Producer is bound to the old client; clear it so next send rebuilds. { let mut producer = self.producer.lock().await; *producer = None; } Ok(()) } async fn validate_tls_files(&self) -> Result<(), TargetError> { // Pulsar only uses CA, no client cert/key. validate_tls_material(&self.args.tls_ca, "", "") } } #[async_trait] impl Target for PulsarTarget where E: PluginEvent, { fn id(&self) -> TargetID { self.id.clone() } async fn is_active(&self) -> Result { self.init_producer().await?; let guard = self.producer.lock().await; let producer = guard .as_ref() .ok_or_else(|| TargetError::Configuration("Pulsar producer not initialized".to_string()))?; producer .check_connection() .await .map_err(|e| TargetError::Network(format!("Pulsar health check failed: {e}")))?; Ok(true) } async fn save(&self, event: Arc>) -> Result<(), TargetError> { let queued = match self.build_queued_payload(&event) { Ok(queued) => queued, Err(err) => { self.delivery_counters.record_final_failure(); return Err(err); } }; if let Some(store) = &self.store { if let Err(e) = persist_queued_payload_to_store(store.as_ref(), &queued) { self.delivery_counters.record_final_failure(); return Err(e); } Ok(()) } else { if let Err(err) = self.send_body(queued.body).await { self.delivery_counters.record_final_failure(); return Err(err); } Ok(()) } } async fn send_raw_from_store(&self, _key: Key, body: Vec, _meta: QueuedPayloadMeta) -> Result<(), TargetError> { self.send_body(body).await } async fn close(&self) -> Result<(), TargetError> { let mut producer = self.producer.lock().await; if let Some(producer) = producer.as_mut() { producer .close() .await .map_err(|e| TargetError::Network(format!("Failed to close Pulsar producer: {e}")))?; } *producer = None; self.clear_cached_client(); self.connected.store(false, Ordering::SeqCst); // If a TLS reload adapter is attached, reset its error tracking // so that a future re-init does not inherit stale failure state. if let Some(adapter) = &self.tls_adapter { *adapter.runtime_state().last_error.write() = None; } info!(target_id = %self.id, "Pulsar target closed"); Ok(()) } fn store(&self) -> Option<&(dyn Store + Send + Sync)> { self.store.as_deref() } fn clone_dyn(&self) -> Box + Send + Sync> { self.clone_box() } async fn init(&self) -> Result<(), TargetError> { if !self.is_enabled() { return Ok(()); } self.init_producer().await } fn is_enabled(&self) -> bool { self.args.enable } fn delivery_snapshot(&self) -> TargetDeliverySnapshot { self.delivery_counters .snapshot(self.store.as_deref().map_or(0, |store| store.len() as u64)) } fn record_final_failure(&self) { self.delivery_counters.record_final_failure(); } } #[cfg(test)] mod tests { use super::*; use crate::target::REDACTED_SECRET; fn base_args() -> PulsarArgs { PulsarArgs { enable: true, broker: "pulsar://127.0.0.1:6650".to_string(), topic: "persistent://public/default/rustfs-events".to_string(), auth_token: String::new(), username: String::new(), password: String::new(), tls_ca: String::new(), tls_allow_insecure: false, tls_hostname_verification: true, queue_dir: String::new(), queue_limit: 0, target_type: TargetType::NotifyEvent, } } #[test] fn debug_redacts_pulsar_secret_fields() { let args = PulsarArgs { auth_token: "pulsar-token".to_string(), password: "pulsar-password".to_string(), ..base_args() }; let rendered = format!("{args:?}"); assert!(!rendered.contains("pulsar-token")); assert!(!rendered.contains("pulsar-password")); assert!(rendered.contains(REDACTED_SECRET)); assert!(rendered.contains("persistent://public/default/rustfs-events")); } #[test] fn validate_pulsar_rejects_mixed_auth_methods() { let args = PulsarArgs { auth_token: "token".to_string(), username: "user".to_string(), password: "pass".to_string(), ..base_args() }; assert!(args.validate().is_err()); } #[test] fn validate_pulsar_rejects_relative_queue_dir() { let args = PulsarArgs { queue_dir: "relative/path".to_string(), ..base_args() }; assert!(args.validate().is_err()); } }