diff --git a/Cargo.lock b/Cargo.lock index f69232ecc..c0e82fdbf 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7200,6 +7200,7 @@ dependencies = [ "query", "rmp-serde", "rust-embed", + "rustfs-event-notifier", "rustfs-obs", "rustls", "rustls-pemfile", @@ -7224,6 +7225,20 @@ dependencies = [ "uuid", ] +[[package]] +name = "rustfs-event-notifier" +version = "0.0.1" +dependencies = [ + "async-trait", + "parking_lot 0.12.3", + "reqwest", + "serde", + "serde_json", + "thiserror 2.0.12", + "tokio", + "wildmatch", +] + [[package]] name = "rustfs-gui" version = "0.0.1" @@ -9317,6 +9332,12 @@ dependencies = [ "rustix 0.38.44", ] +[[package]] +name = "wildmatch" +version = "2.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "68ce1ab1f8c62655ebe1350f589c61e505cf94d385bc6a12899442d9081e71fd" + [[package]] name = "winapi" version = "0.3.9" diff --git a/Cargo.toml b/Cargo.toml index 381305c48..99e5073f3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -15,6 +15,7 @@ members = [ "s3select/api", "s3select/query", "appauth", + "crates/event-notifier", ] resolver = "2" @@ -83,6 +84,7 @@ opentelemetry-stdout = { version = "0.29.0" } opentelemetry-otlp = { version = "0.29" } opentelemetry-prometheus = { version = "0.29.1" } opentelemetry-semantic-conventions = { version = "0.29.0", features = ["semconv_experimental"] } +parking_lot = "0.12.3" pin-project-lite = "0.2" prometheus = "0.14.0" # pin-utils = "0.1.0" @@ -98,6 +100,7 @@ rfd = { version = "0.15.3", default-features = false, features = ["xdg-portal", rmp = "0.8.14" rmp-serde = "1.3.0" rustfs-obs = { path = "crates/obs", version = "0.0.1" } +rustfs-event-notifier = { path = "crates/event-notifier", version = "0.0.1" } rust-embed = "8.7.0" rustls = { version = "0.23.26" } rustls-pki-types = "1.11.0" @@ -141,6 +144,7 @@ uuid = { version = "1.16.0", features = [ "fast-rng", "macro-diagnostics", ] } +wildmatch = "2.4.0" workers = { path = "./common/workers" } [profile.wasm-dev] diff --git a/crates/event-notifier/Cargo.toml b/crates/event-notifier/Cargo.toml new file mode 100644 index 000000000..2e4afc7bc --- /dev/null +++ b/crates/event-notifier/Cargo.toml @@ -0,0 +1,21 @@ +[package] +name = "rustfs-event-notifier" +edition.workspace = true +license.workspace = true +repository.workspace = true +rust-version.workspace = true +version.workspace = true + +[dependencies] +async-trait = { workspace = true } +parking_lot = { workspace = true } +reqwest = { workspace = true } +serde = { workspace = true } +tokio = { workspace = true, features = ["full"] } +thiserror = { workspace = true } +wildmatch = { workspace = true } +serde_json = { workspace = true } + + +[lints] +workspace = true diff --git a/crates/event-notifier/src/error.rs b/crates/event-notifier/src/error.rs new file mode 100644 index 000000000..1945c351c --- /dev/null +++ b/crates/event-notifier/src/error.rs @@ -0,0 +1,25 @@ +// src/error.rs +use thiserror::Error; + +#[derive(Error, Debug)] +pub enum Error { + #[error("target not found: {0}")] + TargetNotFound(String), + + #[error("send failed: {0}")] + SendError(String), + + #[error("target error: {0}")] + TargetError(String), + + #[error("invalid configuration: {0}")] + ConfigError(String), + + #[error("store error: {0}")] + StoreError(String), + + #[error("invalid event name: {0}")] + InvalidEventName(String), // 添加此变体 +} + +pub type Result = std::result::Result; diff --git a/crates/event-notifier/src/event.rs b/crates/event-notifier/src/event.rs new file mode 100644 index 000000000..695f2c3d7 --- /dev/null +++ b/crates/event-notifier/src/event.rs @@ -0,0 +1,53 @@ +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Event { + pub event_version: String, + pub event_source: String, + pub aws_region: String, + pub event_time: String, + pub event_name: String, + pub user_identity: Identity, + pub request_parameters: HashMap, + pub response_elements: HashMap, + pub s3: Metadata, + pub source: Source, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Identity { + pub principal_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Metadata { + pub schema_version: String, + pub configuration_id: String, + pub bucket: Bucket, + pub object: Object, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Bucket { + pub name: String, + pub owner_identity: Identity, + pub arn: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Object { + pub key: String, + pub version_id: Option, + pub sequencer: String, + pub size: Option, + pub etag: Option, + pub content_type: Option, + pub user_metadata: HashMap, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Source { + pub host: String, + pub user_agent: String, +} diff --git a/crates/event-notifier/src/event_name.rs b/crates/event-notifier/src/event_name.rs new file mode 100644 index 000000000..c60cce253 --- /dev/null +++ b/crates/event-notifier/src/event_name.rs @@ -0,0 +1,79 @@ +use crate::error::Error; +use serde::{Deserialize, Serialize}; +use std::str::FromStr; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[repr(u64)] +pub enum EventName { + // 对象操作事件 + ObjectCreatedAll = 1 << 0, + ObjectCreatedPut = 1 << 1, + ObjectCreatedPost = 1 << 2, + ObjectCreatedCopy = 1 << 3, + ObjectCreatedCompleteMultipartUpload = 1 << 4, + + ObjectRemovedAll = 1 << 5, + ObjectRemovedDelete = 1 << 6, + ObjectRemovedDeleteMarkerCreated = 1 << 7, + + ObjectAccessedAll = 1 << 8, + ObjectAccessedGet = 1 << 9, + ObjectAccessedHead = 1 << 10, + + ObjectRestoredAll = 1 << 11, + ObjectRestoredPost = 1 << 12, + ObjectRestoredCompleted = 1 << 13, + + ReplicationAll = 1 << 14, + ReplicationFailed = 1 << 15, + ReplicationComplete = 1 << 16, +} + +impl EventName { + pub fn mask(&self) -> u64 { + *self as u64 + } + + pub fn expand(&self) -> Vec { + match self { + EventName::ObjectCreatedAll => vec![ + EventName::ObjectCreatedPut, + EventName::ObjectCreatedPost, + EventName::ObjectCreatedCopy, + EventName::ObjectCreatedCompleteMultipartUpload, + ], + EventName::ObjectRemovedAll => vec![EventName::ObjectRemovedDelete, EventName::ObjectRemovedDeleteMarkerCreated], + EventName::ObjectAccessedAll => vec![EventName::ObjectAccessedGet, EventName::ObjectAccessedHead], + EventName::ObjectRestoredAll => vec![EventName::ObjectRestoredPost, EventName::ObjectRestoredCompleted], + EventName::ReplicationAll => vec![EventName::ReplicationFailed, EventName::ReplicationComplete], + _ => vec![*self], + } + } +} + +impl FromStr for EventName { + type Err = Error; + + fn from_str(s: &str) -> Result { + match s { + "s3:ObjectCreated:*" => Ok(EventName::ObjectCreatedAll), + "s3:ObjectCreated:Put" => Ok(EventName::ObjectCreatedPut), + "s3:ObjectCreated:Post" => Ok(EventName::ObjectCreatedPost), + "s3:ObjectCreated:Copy" => Ok(EventName::ObjectCreatedCopy), + "s3:ObjectCreated:CompleteMultipartUpload" => Ok(EventName::ObjectCreatedCompleteMultipartUpload), + "s3:ObjectRemoved:*" => Ok(EventName::ObjectRemovedAll), + "s3:ObjectRemoved:Delete" => Ok(EventName::ObjectRemovedDelete), + "s3:ObjectRemoved:DeleteMarkerCreated" => Ok(EventName::ObjectRemovedDeleteMarkerCreated), + "s3:ObjectAccessed:*" => Ok(EventName::ObjectAccessedAll), + "s3:ObjectAccessed:Get" => Ok(EventName::ObjectAccessedGet), + "s3:ObjectAccessed:Head" => Ok(EventName::ObjectAccessedHead), + "s3:ObjectRestored:*" => Ok(EventName::ObjectRestoredAll), + "s3:ObjectRestored:Post" => Ok(EventName::ObjectRestoredPost), + "s3:ObjectRestored:Completed" => Ok(EventName::ObjectRestoredCompleted), + "s3:Replication:*" => Ok(EventName::ReplicationAll), + "s3:Replication:Failed" => Ok(EventName::ReplicationFailed), + "s3:Replication:Complete" => Ok(EventName::ReplicationComplete), + _ => Err(Error::InvalidEventName(format!("Unrecognized event name: {}", s))), + } + } +} diff --git a/crates/event-notifier/src/lib.rs b/crates/event-notifier/src/lib.rs new file mode 100644 index 000000000..fbfc9f863 --- /dev/null +++ b/crates/event-notifier/src/lib.rs @@ -0,0 +1,80 @@ +/// RustFS Event Notifier +/// This crate provides a simple event notification system for RustFS. +/// It allows for the registration of event handlers and the triggering of events. +/// +mod error; +mod event; +mod event_name; +mod notifier; +mod rules; +mod stats; +mod target; +mod target_entry; + +pub fn add(left: u64, right: u64) -> u64 { + left + right +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::event::{Bucket, Event, Identity, Metadata, Object, Source}; + use crate::target::{TargetID, TargetList}; + use std::collections::HashMap; + + #[test] + fn it_works() { + let result = add(2, 2); + assert_eq!(result, 4); + } + + #[tokio::main] + async fn main() { + let target_list = TargetList::new(); + let event = Event { + event_version: "1.0".to_string(), + event_source: "aws:s3".to_string(), + aws_region: "us-west-2".to_string(), + event_time: "2023-10-01T12:00:00Z".to_string(), + event_name: "PutObject".to_string(), + user_identity: Identity { + principal_id: "user123".to_string(), + }, + request_parameters: HashMap::new(), + response_elements: HashMap::new(), + s3: Metadata { + schema_version: "1.0".to_string(), + configuration_id: "config123".to_string(), + bucket: Bucket { + name: "my-bucket".to_string(), + owner_identity: Identity { + principal_id: "owner123".to_string(), + }, + arn: "arn:aws:s3:::my-bucket".to_string(), + }, + object: Object { + key: "my-object.txt".to_string(), + version_id: Some("version123".to_string()), + sequencer: "seq123".to_string(), + size: Some(1024), + etag: Some("etag123".to_string()), + content_type: Some("text/plain".to_string()), + user_metadata: HashMap::new(), + }, + }, + source: Source { + host: "localhost".to_string(), + user_agent: "RustFS/1.0".to_string(), + }, + }; + let target_ids: &[TargetID] = &["".to_string()]; + // 发送事件 + let results = target_list.send(event, &target_ids).await; + println!("result len:{:?}", results.len()); + // 获取统计信息 + let stats = target_list.get_stats(); + for (id, stat) in stats { + println!("Target {}: {} events, {} failed", id, stat.total_events, stat.failed_events); + } + } +} diff --git a/crates/event-notifier/src/notifier.rs b/crates/event-notifier/src/notifier.rs new file mode 100644 index 000000000..7b81a24d2 --- /dev/null +++ b/crates/event-notifier/src/notifier.rs @@ -0,0 +1,47 @@ +// src/notifier.rs +use crate::error::Error; +use crate::event::Event; +use crate::rules::{RulesMap, TargetIDSet}; +use crate::target::TargetList; +use std::collections::HashMap; +use std::sync::{Arc, RwLock}; +use tokio::sync::broadcast; + +pub struct EventNotifier { + target_list: Arc, + rules: RwLock>, + tx: broadcast::Sender, +} + +impl EventNotifier { + pub fn new(capacity: usize) -> Self { + let (tx, _) = broadcast::channel(capacity); + Self { + target_list: Arc::new(TargetList::new()), + rules: RwLock::new(HashMap::new()), + tx, + } + } + + pub async fn send(&self, event: Event) -> Result<(), Error> { + let rules = self + .rules + .read() + .map_err(|_| Error::StoreError("Failed to read rules".to_string()))?; + // 检查规则匹配 + let target_ids = if let Some(rules_map) = rules.get(&event.s3.bucket.name) { + rules_map.match_targets(event.s3.event_name, &event.s3.object.key) + } else { + TargetIDSet::new() + }; + + // 发送事件 + if !target_ids.is_empty() { + self.target_list.send(event.clone(), &target_ids).await; + } + + // 广播给监听者 + let _ = self.tx.send(event); + Ok(()) + } +} diff --git a/crates/event-notifier/src/rules.rs b/crates/event-notifier/src/rules.rs new file mode 100644 index 000000000..671bcfb53 --- /dev/null +++ b/crates/event-notifier/src/rules.rs @@ -0,0 +1,191 @@ +// src/rules.rs +use crate::event_name::EventName; +use crate::target::TargetID; +use std::collections::{HashMap, HashSet}; +use wildmatch::WildMatch; + +/// TargetIDSet represents a set of target IDs +#[derive(Debug, Clone, Default)] +pub struct TargetIDSet(HashSet); + +impl TargetIDSet { + pub fn new() -> Self { + Self(HashSet::new()) + } + + pub fn insert(&mut self, target_id: TargetID) -> bool { + self.0.insert(target_id) + } + + pub fn contains(&self, target_id: &TargetID) -> bool { + self.0.contains(target_id) + } + + pub fn remove(&mut self, target_id: &TargetID) -> bool { + self.0.remove(target_id) + } + + pub fn len(&self) -> usize { + self.0.len() + } + + pub fn is_empty(&self) -> bool { + self.0.is_empty() + } + + pub fn clear(&mut self) { + self.0.clear() + } + + pub fn iter(&self) -> impl Iterator { + self.0.iter() + } + + pub fn extend>(&mut self, iter: I) { + self.0.extend(iter) + } + + pub fn union(&self, other: &TargetIDSet) -> TargetIDSet { + let mut result = self.clone(); + result.extend(other.0.iter().cloned()); + result + } + + pub fn difference(&self, other: &TargetIDSet) -> TargetIDSet { + let mut result = TargetIDSet::new(); + for id in self.0.iter() { + if !other.contains(id) { + result.insert(id.clone()); + } + } + result + } +} + +impl FromIterator for TargetIDSet { + fn from_iter>(iter: I) -> Self { + let mut set = TargetIDSet::new(); + set.extend(iter); + set + } +} + +#[derive(Debug, Clone)] +pub struct Rules { + patterns: HashMap, +} + +impl Rules { + pub fn new() -> Self { + Self { + patterns: HashMap::new(), + } + } + + pub fn add(&mut self, pattern: String, target_id: TargetID) { + let entry = self.patterns.entry(pattern).or_insert_with(TargetIDSet::new); + entry.insert(target_id); + } + + pub fn match_object(&self, object_name: &str) -> bool { + for pattern in self.patterns.keys() { + if WildMatch::new(pattern).matches(object_name) { + return true; + } + } + false + } + + pub fn match_targets(&self, object_name: &str) -> TargetIDSet { + let mut target_ids = TargetIDSet::new(); + + for (pattern, ids) in &self.patterns { + if WildMatch::new(pattern).matches(object_name) { + target_ids.extend(ids.iter().cloned()); + } + } + + target_ids + } +} + +#[derive(Debug, Default)] +pub struct RulesMap { + rules: HashMap, +} + +impl RulesMap { + pub fn new() -> Self { + Self { rules: HashMap::new() } + } + + pub fn add(&mut self, events: &[EventName], pattern: String, target_id: TargetID) { + for &event in events { + for expanded_event in event.expand() { + let rules = self.rules.entry(expanded_event).or_insert_with(Rules::new); + rules.add(pattern.clone(), target_id.clone()); + } + } + } + + pub fn match_simple(&self, event: EventName, object_name: &str) -> bool { + match self.rules.get(&event) { + Some(rules) => rules.match_object(object_name), + None => false, + } + } + + pub fn match_targets(&self, event: EventName, object_name: &str) -> TargetIDSet { + match self.rules.get(&event) { + Some(rules) => rules.match_targets(object_name), + None => TargetIDSet::new(), + } + } +} + +impl Clone for RulesMap { + fn clone(&self) -> Self { + Self { + rules: self.rules.clone(), + } + } +} + +impl Default for Rules { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::event_name::EventName; + + #[test] + fn test_rules() { + let mut rules = Rules::new(); + rules.add("*.txt".to_string(), "target1".to_string()); + rules.add("*.jpg".to_string(), "target2".to_string()); + + assert!(rules.match_object("file.txt")); + assert!(!rules.match_object("file.pdf")); + assert!(rules.match_object("image.jpg")); + assert!(!rules.match_object("image.png")); + + let target_ids = rules.match_targets("file.txt"); + assert_eq!(target_ids.len(), 1); + assert!(target_ids.iter().any(|id| id == "target1")); + + let mut rules_map = RulesMap::new(); + let target_id = "target1".to_string(); + // 添加规则 + rules_map.add(&[EventName::ObjectCreatedAll], String::from("images/*"), target_id); + + // 匹配对象 + let matches = rules_map.match_simple(EventName::ObjectCreatedPut, "images/photo.jpg"); // returns true + + // 获取匹配的目标 + let target_ids = rules_map.match_targets(EventName::ObjectCreatedPut, "images/photo.jpg"); + } +} diff --git a/crates/event-notifier/src/stats.rs b/crates/event-notifier/src/stats.rs new file mode 100644 index 000000000..570bbeca5 --- /dev/null +++ b/crates/event-notifier/src/stats.rs @@ -0,0 +1,81 @@ +// src/stats.rs +use crate::target::TargetID; +use parking_lot::RwLock; +use std::collections::HashMap; +use std::sync::atomic::{AtomicI64, Ordering}; +use std::sync::Arc; + +#[derive(Debug, Default)] +pub struct TargetStat { + pub current_send_calls: AtomicI64, + pub total_events: AtomicI64, + pub failed_events: AtomicI64, + pub current_queue: AtomicI64, +} + +impl TargetStat { + pub fn inc_send_calls(&self) { + self.current_send_calls.fetch_add(1, Ordering::Relaxed); + } + + pub fn dec_send_calls(&self) { + self.current_send_calls.fetch_sub(1, Ordering::Relaxed); + } + + pub fn inc_total_events(&self) { + self.total_events.fetch_add(1, Ordering::Relaxed); + } + + pub fn inc_failed_events(&self) { + self.failed_events.fetch_add(1, Ordering::Relaxed); + } + + pub fn set_queue_size(&self, size: i64) { + self.current_queue.store(size, Ordering::Relaxed); + } +} + +#[derive(Default)] +pub struct TargetStats { + stats: RwLock>>, +} + +impl TargetStats { + pub fn new() -> Self { + Self::default() + } + + pub fn get_or_create(&self, id: &TargetID) -> Arc { + let mut stats = self.stats.write(); + stats + .entry(id.clone()) + .or_insert_with(|| Arc::new(TargetStat::default())) + .clone() + } + + pub fn get_stats(&self) -> HashMap { + self.stats + .read() + .iter() + .map(|(id, stat)| { + ( + id.clone(), + TargetSnapshot { + current_send_calls: stat.current_send_calls.load(Ordering::Relaxed), + total_events: stat.total_events.load(Ordering::Relaxed), + failed_events: stat.failed_events.load(Ordering::Relaxed), + current_queue: stat.current_queue.load(Ordering::Relaxed), + }, + ) + }) + .collect() + } +} + +#[derive(Debug, Clone)] +pub struct TargetSnapshot { + pub current_send_calls: i64, + pub total_events: i64, + pub failed_events: i64, + pub current_queue: i64, +} diff --git a/crates/event-notifier/src/target.rs b/crates/event-notifier/src/target.rs new file mode 100644 index 000000000..68b575b74 --- /dev/null +++ b/crates/event-notifier/src/target.rs @@ -0,0 +1,66 @@ +// src/target.rs +use crate::error::{Error, Result}; +use crate::event::Event; +use crate::stats::{TargetSnapshot, TargetStats}; +use async_trait::async_trait; +use std::collections::HashMap; +use std::sync::{Arc, RwLock}; + +pub type TargetID = String; + +#[async_trait] +pub trait Target: Send + Sync { + fn id(&self) -> TargetID; + async fn send(&self, event: Event) -> Result<()>; + async fn is_active(&self) -> bool; + async fn close(&self) -> Result<()>; +} + +pub struct TargetList { + targets: Arc>>>, + stats: Arc, +} + +impl TargetList { + pub fn new() -> Self { + Self { + targets: Arc::new(RwLock::new(HashMap::new())), + stats: Arc::default(), + } + } + + pub async fn send(&self, event: Event, target_ids: &[TargetID]) -> Vec> { + let mut results = Vec::with_capacity(target_ids.len()); + let targets = self.targets.read().unwrap(); + + for id in target_ids { + let target = match targets.get(id) { + Some(t) => t, + None => { + results.push(Err(Error::TargetNotFound(id.to_string()))); + continue; + } + }; + + let stats = self.stats.get_or_create(id); + stats.inc_send_calls(); + + let result = target.send(event.clone()).await; + + stats.dec_send_calls(); + stats.inc_total_events(); + + if result.is_err() { + stats.inc_failed_events(); + } + + results.push(result); + } + + results + } + + pub fn get_stats(&self) -> HashMap { + self.stats.get_stats() + } +} diff --git a/crates/event-notifier/src/target_entry/mod.rs b/crates/event-notifier/src/target_entry/mod.rs new file mode 100644 index 000000000..a2a1c2754 --- /dev/null +++ b/crates/event-notifier/src/target_entry/mod.rs @@ -0,0 +1 @@ +mod webhook; diff --git a/crates/event-notifier/src/target_entry/webhook.rs b/crates/event-notifier/src/target_entry/webhook.rs new file mode 100644 index 000000000..1736ad36b --- /dev/null +++ b/crates/event-notifier/src/target_entry/webhook.rs @@ -0,0 +1,161 @@ +use crate::event::Event; +use crate::target::{Target, TargetID}; +use async_trait::async_trait; +use reqwest::{header, Client}; +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; +use std::sync::Arc; +use std::time::Duration; +use thiserror::Error; +use tokio::io::AsyncWriteExt; +use tokio::sync::Mutex; + +#[derive(Error, Debug)] +pub enum WebhookError { + #[error("HTTP client error: {0}")] + HttpClient(#[from] reqwest::Error), + #[error("JSON serialization error: {0}")] + Json(#[from] serde_json::Error), + #[error("Store error: {0}")] + Store(#[from] Box), + #[error("Webhook request failed after retries")] + RequestFailed, +} + +type Result = std::result::Result; + +/// Webhook 目标配置 + +/// Webhook 目标配置 +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WebhookConfig { + /// Webhook endpoint URL + pub endpoint: String, + /// 认证令牌 + pub auth_token: Option, + /// 自定义请求头 + pub custom_headers: Option>, + /// 重试次数 + pub max_retries: u32, + /// 连接超时时间 (秒) + pub timeout: u64, +} + +/// Webhook Target 实现 +pub struct Webhook { + /// 目标 ID + id: TargetID, + /// HTTP 客户端 + client: Client, + /// 配置信息 + config: WebhookConfig, + /// 事件存储 + store: Arc>, +} + +#[async_trait] +impl Target for Webhook { + /// 返回目标 ID + fn id(&self) -> &TargetID { + &self.id + } + + /// 检查 Webhook 是否可用 + async fn is_active(&self) -> Result { + // 发送测试请求验证连接 + let resp = self.client.get(&self.config.endpoint).send().await?; + + Ok(resp.status().is_success()) + } + + /// 保存事件到存储并异步发送 + async fn save(&self, event: Event) -> Result<()> { + // 序列化事件 + let event_json = serde_json::to_string(&event)?; + + // 存储事件 + self.store.lock().await.save(event_json)?; + + // 异步发送 + tokio::spawn(self.send_event(event)); + + Ok(()) + } + + /// 从存储发送事件 + async fn send(&self, key: StoreKey) -> Result<()> { + // 从存储获取事件 + let event_json = self.store.lock().await.get(&key)?; + let event: Event = serde_json::from_str(&event_json)?; + + // 发送事件到 webhook + self.send_event(event).await?; + + // 发送成功后删除 + self.store.lock().await.delete(&key)?; + + Ok(()) + } + + /// 返回事件存储 + fn store(&self) -> Arc> { + self.store.clone() + } + + /// 关闭 Target + async fn close(&self) -> Result<()> { + // 等待所有事件发送完成 + self.store.lock().await.flush()?; + Ok(()) + } +} + +impl Webhook { + /// 创建新的 Webhook Target + pub fn new(id: TargetID, config: WebhookConfig, store: Arc>) -> Result { + // 构建 HTTP 客户端 + let client = Client::builder().timeout(Duration::from_secs(config.timeout)).build()?; + + Ok(Self { + id, + client, + config, + store, + }) + } + + /// 发送事件到 Webhook endpoint + async fn send_event(&self, event: Event) -> Result<()> { + let mut retries = 0; + loop { + // 构建请求 + let mut req = self.client.post(&self.config.endpoint).json(&event); + + // 添加认证头 + if let Some(token) = &self.config.auth_token { + req = req.header(header::AUTHORIZATION, format!("Bearer {}", token)); + } + + // 添加自定义头 + if let Some(headers) = &self.config.custom_headers { + for (key, value) in headers { + req = req.header(key, value); + } + } + + // 发送请求 + match req.send().await { + Ok(resp) if resp.status().is_success() => { + return Ok(()); + } + _ if retries < self.config.max_retries => { + retries += 1; + tokio::time::sleep(Duration::from_secs(1)).await; + continue; + } + Err(e) => return Err(WebhookError::HttpClient(e)), + _ => return Err(WebhookError::RequestFailed), + } + } + } +} diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index ff8b5971e..e2d088156 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -50,6 +50,7 @@ pin-project-lite.workspace = true protos.workspace = true query = { path = "../s3select/query" } rmp-serde.workspace = true +rustfs-event-notifier = { workspace = true } rustfs-obs = { workspace = true } rustls.workspace = true rustls-pemfile.workspace = true