mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-06 05:17:42 +00:00
c5df1f92c2
* improve code for notify * improve code for logger and fix typo (#272) * Add GNU to build.yml (#275) * fix unzip error * fix url change error fix url change error * Simplify user experience and integrate console and endpoint Simplify user experience and integrate console and endpoint * Add gnu to build.yml * upgrade version * feat: add `cargo clippy --fix --allow-dirty` to pre-commit command (#282) Resolves #277 - Add --fix flag to automatically fix clippy warnings - Add --allow-dirty flag to run on dirty Git trees - Improves code quality in pre-commit workflow * fix: the issue where preview fails when the path length exceeds 255 characters (#280) * fix * fix: improve Windows build support and CI/CD workflow (#283) - Fix Windows zip command issue by using PowerShell Compress-Archive - Add Windows support for OSS upload with ossutil - Replace Chinese comments with English in build.yml - Fix bash syntax error in package_zip function - Improve code formatting and consistency - Update various configuration files for better cross-platform support Resolves Windows build failures in GitHub Actions. * fix: update link in README.md leading to a 404 error (#285) * add rustfs.spec for rustfs (#103) add support on loongarch64 * improve cargo.lock * build(deps): bump the dependencies group with 5 updates (#289) Bumps the dependencies group with 5 updates: | Package | From | To | | --- | --- | --- | | [hyper-util](https://github.com/hyperium/hyper-util) | `0.1.15` | `0.1.16` | | [rand](https://github.com/rust-random/rand) | `0.9.1` | `0.9.2` | | [serde_json](https://github.com/serde-rs/json) | `1.0.140` | `1.0.141` | | [strum](https://github.com/Peternator7/strum) | `0.27.1` | `0.27.2` | | [sysinfo](https://github.com/GuillaumeGomez/sysinfo) | `0.36.0` | `0.36.1` | Updates `hyper-util` from 0.1.15 to 0.1.16 - [Release notes](https://github.com/hyperium/hyper-util/releases) - [Changelog](https://github.com/hyperium/hyper-util/blob/master/CHANGELOG.md) - [Commits](https://github.com/hyperium/hyper-util/compare/v0.1.15...v0.1.16) Updates `rand` from 0.9.1 to 0.9.2 - [Release notes](https://github.com/rust-random/rand/releases) - [Changelog](https://github.com/rust-random/rand/blob/master/CHANGELOG.md) - [Commits](https://github.com/rust-random/rand/compare/rand_core-0.9.1...rand_core-0.9.2) Updates `serde_json` from 1.0.140 to 1.0.141 - [Release notes](https://github.com/serde-rs/json/releases) - [Commits](https://github.com/serde-rs/json/compare/v1.0.140...v1.0.141) Updates `strum` from 0.27.1 to 0.27.2 - [Release notes](https://github.com/Peternator7/strum/releases) - [Changelog](https://github.com/Peternator7/strum/blob/master/CHANGELOG.md) - [Commits](https://github.com/Peternator7/strum/compare/v0.27.1...v0.27.2) Updates `sysinfo` from 0.36.0 to 0.36.1 - [Changelog](https://github.com/GuillaumeGomez/sysinfo/blob/master/CHANGELOG.md) - [Commits](https://github.com/GuillaumeGomez/sysinfo/compare/v0.36.0...v0.36.1) --- updated-dependencies: - dependency-name: hyper-util dependency-version: 0.1.16 dependency-type: direct:production update-type: version-update:semver-patch dependency-group: dependencies - dependency-name: rand dependency-version: 0.9.2 dependency-type: direct:production update-type: version-update:semver-patch dependency-group: dependencies - dependency-name: serde_json dependency-version: 1.0.141 dependency-type: direct:production update-type: version-update:semver-patch dependency-group: dependencies - dependency-name: strum dependency-version: 0.27.2 dependency-type: direct:production update-type: version-update:semver-patch dependency-group: dependencies - dependency-name: sysinfo dependency-version: 0.36.1 dependency-type: direct:production update-type: version-update:semver-patch dependency-group: dependencies ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> * improve code for logger * improve * upgrade * refactor: 优化构建工作流,统一 latest 文件处理和简化制品上传 (#293) * Refactor: DatabaseManagerSystem as global Signed-off-by: junxiang Mu <1948535941@qq.com> * fix: fmt Signed-off-by: junxiang Mu <1948535941@qq.com> * Test: add e2e_test for s3select Signed-off-by: junxiang Mu <1948535941@qq.com> * Test: add test script for e2e Signed-off-by: junxiang Mu <1948535941@qq.com> * improve code for registry and intergation * improve code for registry `create_targets_from_config` * fix * Feature up/ilm (#305) * fix * fix * fix * fix delete-marker expiration. add api_restore. * fix * time retry object upload * lock file * make fmt * fix * restore object * fix * fix * serde-rs-xml -> quick-xml * fix * checksum * fix * fix * fix * fix * fix * fix * fix * transfer lang to english * upgrade clap version from 4.5.41 to 4.5.42 * refactor: replace `lazy_static` with `LazyLock` * add router * fix: modify comment * improve code * fix typos * fix * fix: modify name and fmt * improve code for registry * fix test --------- Signed-off-by: dependabot[bot] <support@github.com> Signed-off-by: junxiang Mu <1948535941@qq.com> Co-authored-by: loverustfs <155562731+loverustfs@users.noreply.github.com> Co-authored-by: 安正超 <anzhengchao@gmail.com> Co-authored-by: shiro.lee <69624924+shiroleeee@users.noreply.github.com> Co-authored-by: Marco Orlandin <mipnamic@mipnamic.net> Co-authored-by: zhangwenlong <zhangwenlong@loongson.cn> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> Co-authored-by: junxiang Mu <1948535941@qq.com> Co-authored-by: likewu <likewu@126.com>
186 lines
6.7 KiB
Rust
186 lines
6.7 KiB
Rust
// 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 rustfs_config::notify::{
|
|
DEFAULT_LIMIT, DEFAULT_TARGET, ENABLE_KEY, ENABLE_ON, MQTT_BROKER, MQTT_PASSWORD, MQTT_QOS, MQTT_QUEUE_DIR, MQTT_QUEUE_LIMIT,
|
|
MQTT_TOPIC, MQTT_USERNAME, NOTIFY_MQTT_SUB_SYS, NOTIFY_WEBHOOK_SUB_SYS, WEBHOOK_AUTH_TOKEN, WEBHOOK_ENDPOINT,
|
|
WEBHOOK_QUEUE_DIR, WEBHOOK_QUEUE_LIMIT,
|
|
};
|
|
use rustfs_ecstore::config::{Config, KV, KVS};
|
|
use rustfs_notify::arn::TargetID;
|
|
use rustfs_notify::{BucketNotificationConfig, Event, EventName, LogLevel, NotificationError, init_logger};
|
|
use rustfs_notify::{initialize, notification_system};
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
use tracing::info;
|
|
|
|
#[tokio::main]
|
|
async fn main() -> Result<(), NotificationError> {
|
|
init_logger(LogLevel::Debug);
|
|
|
|
let system = match notification_system() {
|
|
Some(sys) => sys,
|
|
None => {
|
|
let config = Config::new();
|
|
initialize(config).await?;
|
|
notification_system().expect("Failed to initialize notification system")
|
|
}
|
|
};
|
|
|
|
// --- Initial configuration (Webhook and MQTT) ---
|
|
let mut config = Config::new();
|
|
let current_root = rustfs_utils::dirs::get_project_root().expect("failed to get project root");
|
|
println!("Current project root: {}", current_root.display());
|
|
|
|
let webhook_kvs_vec = vec![
|
|
KV {
|
|
key: ENABLE_KEY.to_string(),
|
|
value: ENABLE_ON.to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: WEBHOOK_ENDPOINT.to_string(),
|
|
value: "http://127.0.0.1:3020/webhook".to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: WEBHOOK_AUTH_TOKEN.to_string(),
|
|
value: "secret-token".to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: WEBHOOK_QUEUE_DIR.to_string(),
|
|
value: current_root
|
|
.clone()
|
|
.join("../../deploy/logs/notify/webhook")
|
|
.to_str()
|
|
.unwrap()
|
|
.to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: WEBHOOK_QUEUE_LIMIT.to_string(),
|
|
value: DEFAULT_LIMIT.to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
];
|
|
let webhook_kvs = KVS(webhook_kvs_vec);
|
|
|
|
let mut webhook_targets = std::collections::HashMap::new();
|
|
webhook_targets.insert(DEFAULT_TARGET.to_string(), webhook_kvs);
|
|
config.0.insert(NOTIFY_WEBHOOK_SUB_SYS.to_string(), webhook_targets);
|
|
|
|
// MQTT target configuration
|
|
let mqtt_kvs_vec = vec![
|
|
KV {
|
|
key: ENABLE_KEY.to_string(),
|
|
value: ENABLE_ON.to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: MQTT_BROKER.to_string(),
|
|
value: "mqtt://localhost:1883".to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: MQTT_TOPIC.to_string(),
|
|
value: "rustfs/events".to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: MQTT_QOS.to_string(),
|
|
value: "1".to_string(), // AtLeastOnce
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: MQTT_USERNAME.to_string(),
|
|
value: "test".to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: MQTT_PASSWORD.to_string(),
|
|
value: "123456".to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: MQTT_QUEUE_DIR.to_string(),
|
|
value: current_root
|
|
.join("../../deploy/logs/notify/mqtt")
|
|
.to_str()
|
|
.unwrap()
|
|
.to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
KV {
|
|
key: MQTT_QUEUE_LIMIT.to_string(),
|
|
value: DEFAULT_LIMIT.to_string(),
|
|
hidden_if_empty: false,
|
|
},
|
|
];
|
|
|
|
let mqtt_kvs = KVS(mqtt_kvs_vec);
|
|
let mut mqtt_targets = std::collections::HashMap::new();
|
|
mqtt_targets.insert(DEFAULT_TARGET.to_string(), mqtt_kvs);
|
|
config.0.insert(NOTIFY_MQTT_SUB_SYS.to_string(), mqtt_targets);
|
|
|
|
// Load the configuration and initialize the system
|
|
*system.config.write().await = config;
|
|
system.init().await?;
|
|
info!("✅ System initialized with Webhook and MQTT targets.");
|
|
|
|
// --- Query the currently active Target ---
|
|
let active_targets = system.get_active_targets().await;
|
|
info!("\n---> Currently active targets: {:?}", active_targets);
|
|
assert_eq!(active_targets.len(), 2);
|
|
|
|
tokio::time::sleep(Duration::from_secs(1)).await;
|
|
|
|
// --- Exactly delete a Target (e.g. MQTT) ---
|
|
info!("\n---> Removing MQTT target...");
|
|
let mqtt_target_id = TargetID::new(DEFAULT_TARGET.to_string(), "mqtt".to_string());
|
|
system.remove_target(&mqtt_target_id, NOTIFY_MQTT_SUB_SYS).await?;
|
|
info!("✅ MQTT target removed.");
|
|
|
|
// --- Query the activity's Target again ---
|
|
let active_targets_after_removal = system.get_active_targets().await;
|
|
info!("\n---> Active targets after removal: {:?}", active_targets_after_removal);
|
|
assert_eq!(active_targets_after_removal.len(), 1);
|
|
assert_eq!(active_targets_after_removal[0].id, DEFAULT_TARGET.to_string());
|
|
|
|
// --- Send events for verification ---
|
|
// Configure a rule to point to the Webhook and deleted MQTT
|
|
let mut bucket_config = BucketNotificationConfig::new("us-east-1");
|
|
bucket_config.add_rule(
|
|
&[EventName::ObjectCreatedPut],
|
|
"*".to_string(),
|
|
TargetID::new(DEFAULT_TARGET.to_string(), "webhook".to_string()),
|
|
);
|
|
bucket_config.add_rule(
|
|
&[EventName::ObjectCreatedPut],
|
|
"*".to_string(),
|
|
TargetID::new(DEFAULT_TARGET.to_string(), "mqtt".to_string()), // This rule will match, but the Target cannot be found
|
|
);
|
|
system.load_bucket_notification_config("my-bucket", &bucket_config).await?;
|
|
|
|
info!("\n---> Sending an event...");
|
|
let event = Arc::new(Event::new_test_event("my-bucket", "document.pdf", EventName::ObjectCreatedPut));
|
|
system.send_event(event).await;
|
|
info!("✅ Event sent. Only the Webhook target should receive it. Check logs for warnings about the missing MQTT target.");
|
|
|
|
tokio::time::sleep(Duration::from_secs(2)).await;
|
|
|
|
info!("\nDemo completed successfully");
|
|
Ok(())
|
|
}
|