From 434b71e569eac93121dcb82c40b3dcdfa98eeb63 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=AE=89=E6=AD=A3=E8=B6=85?= Date: Fri, 12 Jun 2026 15:54:42 +0800 Subject: [PATCH] feat(extension): add admin views (#3386) --- Cargo.lock | 1 + rustfs/Cargo.toml | 1 + rustfs/src/admin/handlers/extensions.rs | 312 ++++++++++++++++++ rustfs/src/admin/handlers/mod.rs | 3 + .../src/admin/handlers/plugins_instances.rs | 15 +- rustfs/src/admin/mod.rs | 6 +- rustfs/src/admin/route_registration_test.rs | 4 + 7 files changed, 334 insertions(+), 8 deletions(-) create mode 100644 rustfs/src/admin/handlers/extensions.rs diff --git a/Cargo.lock b/Cargo.lock index c5b9eafc6..a9cbbb7ce 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9074,6 +9074,7 @@ dependencies = [ "rustfs-crypto", "rustfs-data-usage", "rustfs-ecstore", + "rustfs-extension-schema", "rustfs-filemeta", "rustfs-heal", "rustfs-iam", diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 3b8f19e27..c36cd6a13 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -86,6 +86,7 @@ rustfs-security-governance = { workspace = true } rustfs-data-usage = { workspace = true } rustfs-s3select-api = { workspace = true } rustfs-s3select-query = { workspace = true } +rustfs-extension-schema = { workspace = true } rustfs-targets = { workspace = true } rustfs-storage-api = { workspace = true } rustfs-tls-runtime = { workspace = true } diff --git a/rustfs/src/admin/handlers/extensions.rs b/rustfs/src/admin/handlers/extensions.rs new file mode 100644 index 000000000..ba6e758c4 --- /dev/null +++ b/rustfs/src/admin/handlers/extensions.rs @@ -0,0 +1,312 @@ +// 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::admin::{ + auth::validate_admin_request, + handlers::plugins_instances, + plugin_contract::{ + PluginContractDomain, PluginInstanceDiagnosticCode, PluginInstanceDiagnosticCount, PluginInstanceEntry, + PluginInstanceSource, PluginOperationalStateContract, + }, + router::{AdminOperation, Operation, S3Router}, +}; +use crate::auth::{check_key_valid, get_session_token}; +use crate::server::{ADMIN_PREFIX, RemoteAddr}; +use http::{HeaderMap, HeaderValue, StatusCode}; +use hyper::Method; +use matchit::Params; +use rustfs_extension_schema::{ExtensionKind, ExtensionSchema}; +use rustfs_policy::policy::action::{Action, AdminAction}; +use rustfs_targets::{ + builtin_target_extension_schemas, catalog::example_external_webhook_plugin, target_marketplace_extension_schema, +}; +use s3s::header::CONTENT_TYPE; +use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; +use serde::Serialize; +use std::collections::HashMap; + +pub fn register_extension_route(r: &mut S3Router) -> std::io::Result<()> { + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/v4/extensions/catalog").as_str(), + AdminOperation(&GetExtensionCatalogHandler {}), + )?; + r.insert( + Method::GET, + format!("{}{}", ADMIN_PREFIX, "/v4/extensions/instances").as_str(), + AdminOperation(&ListExtensionInstancesHandler {}), + )?; + + Ok(()) +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub(crate) struct ExtensionCatalogResponse { + pub extensions: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub(crate) struct ExtensionInstanceEntry { + pub id: String, + pub extension_id: String, + pub kind: ExtensionKind, + pub domain: PluginContractDomain, + pub subsystem: String, + pub account_id: String, + pub service: String, + pub status: String, + pub source: PluginInstanceSource, + pub enabled: bool, + pub config: HashMap, + #[serde(skip_serializing_if = "Option::is_none")] + pub operational_state: Option, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub diagnostic_codes: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub(crate) struct ExtensionInstancesResponse { + pub instances: Vec, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub diagnostic_counts: Vec, + pub truncated: bool, + pub next_marker: Option, +} + +fn build_extension_catalog_response() -> ExtensionCatalogResponse { + let mut extensions = builtin_target_extension_schemas(); + let example = example_external_webhook_plugin(); + extensions.push(target_marketplace_extension_schema(&example.manifest)); + extensions.sort_by(|a, b| a.extension_id.cmp(&b.extension_id)); + + ExtensionCatalogResponse { extensions } +} + +fn map_extension_instance(instance: PluginInstanceEntry) -> ExtensionInstanceEntry { + ExtensionInstanceEntry { + id: instance.id, + extension_id: instance.plugin_id, + kind: ExtensionKind::TargetPlugin, + domain: instance.domain, + subsystem: instance.subsystem, + account_id: instance.account_id, + service: instance.service, + status: instance.status, + source: instance.source, + enabled: instance.enabled, + config: instance.config, + operational_state: instance.operational_state, + diagnostic_codes: instance.diagnostic_codes, + } +} + +async fn authorize_extension_catalog_request(req: &S3Request) -> S3Result<()> { + let Some(input_cred) = &req.credentials else { + return Err(s3_error!(InvalidRequest, "authentication required")); + }; + + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; + + validate_admin_request( + &req.headers, + &cred, + owner, + false, + vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)], + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await +} + +async fn authorize_extension_instance_request(req: &S3Request) -> S3Result<()> { + let Some(input_cred) = &req.credentials else { + return Err(s3_error!(InvalidRequest, "authentication required")); + }; + + let (cred, owner) = + check_key_valid(get_session_token(&req.uri, &req.headers).unwrap_or_default(), &input_cred.access_key).await?; + + validate_admin_request( + &req.headers, + &cred, + owner, + false, + vec![Action::AdminAction(AdminAction::GetBucketTargetAction)], + req.extensions.get::>().and_then(|opt| opt.map(|a| a.0)), + ) + .await +} + +fn build_json_response( + status: StatusCode, + body: &impl Serialize, + request_id: Option<&HeaderValue>, +) -> S3Result> { + let data = serde_json::to_vec(body).map_err(|e| s3_error!(InternalError, "failed to serialize response: {}", e))?; + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, HeaderValue::from_static("application/json")); + if let Some(value) = request_id { + header.insert("x-request-id", value.clone()); + } + Ok(S3Response::with_headers((status, Body::from(data)), header)) +} + +pub struct GetExtensionCatalogHandler {} + +#[async_trait::async_trait] +impl Operation for GetExtensionCatalogHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + authorize_extension_catalog_request(&req).await?; + build_json_response(StatusCode::OK, &build_extension_catalog_response(), req.headers.get("x-request-id")) + } +} + +pub struct ListExtensionInstancesHandler {} + +#[async_trait::async_trait] +impl Operation for ListExtensionInstancesHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + authorize_extension_instance_request(&req).await?; + let filters = plugins_instances::extract_plugin_instance_filters(&req)?; + let instances = plugins_instances::filter_plugin_instances(plugins_instances::collect_all_instances().await?, &filters); + let diagnostic_counts = plugins_instances::collect_diagnostic_counts(&instances); + let (instances, truncated, next_marker) = plugins_instances::paginate_plugin_instances(instances, &filters)?; + let response = ExtensionInstancesResponse { + instances: instances.into_iter().map(map_extension_instance).collect(), + diagnostic_counts, + truncated, + next_marker, + }; + build_json_response(StatusCode::OK, &response, req.headers.get("x-request-id")) + } +} + +#[cfg(test)] +mod tests { + use super::{build_extension_catalog_response, map_extension_instance}; + use crate::admin::plugin_contract::{ + PluginContractDomain, PluginEnableState, PluginInstallState, PluginInstanceEntry, PluginInstanceSource, + PluginOperationalRuntimeState, PluginOperationalStateContract, + }; + use rustfs_extension_schema::{ExtensionKind, ExtensionRuntimeBoundary, validate_extension_schemas}; + use std::collections::HashMap; + + #[test] + fn extension_handlers_require_admin_authorization_contract() { + let src = include_str!("extensions.rs"); + let catalog_block = extract_block_between_markers( + src, + "impl Operation for GetExtensionCatalogHandler", + "pub struct ListExtensionInstancesHandler", + ); + let instances_block = + extract_block_between_markers(src, "impl Operation for ListExtensionInstancesHandler", "#[cfg(test)]"); + + assert!( + catalog_block.contains("authorize_extension_catalog_request(&req).await?;"), + "extension catalog should require admin authorization" + ); + assert!( + instances_block.contains("authorize_extension_instance_request(&req).await?;"), + "extension instances should require admin authorization" + ); + + let catalog_auth = extract_block_between_markers( + src, + "async fn authorize_extension_catalog_request", + "async fn authorize_extension_instance_request", + ); + assert!( + catalog_auth.contains("AdminAction::ServerInfoAdminAction"), + "extension catalog should require server info admin permission" + ); + let instance_auth = + extract_block_between_markers(src, "async fn authorize_extension_instance_request", "fn build_json_response"); + assert!( + instance_auth.contains("AdminAction::GetBucketTargetAction"), + "extension instances should require target read permission" + ); + } + + #[test] + fn extension_catalog_exposes_valid_target_plugin_schemas() { + let response = build_extension_catalog_response(); + let webhook = response + .extensions + .iter() + .find(|schema| schema.extension_id == "builtin:webhook") + .expect("builtin webhook extension should be present"); + + assert_eq!(webhook.kind, ExtensionKind::TargetPlugin); + assert_eq!(webhook.runtime.boundary, ExtensionRuntimeBoundary::Builtin); + assert!(!webhook.disabled_by_default); + assert!( + webhook + .capabilities + .iter() + .any(|capability| capability.as_str() == "target.audit.v1") + ); + assert!( + webhook + .capabilities + .iter() + .any(|capability| capability.as_str() == "target.notify.v1") + ); + assert!(validate_extension_schemas(&response.extensions).is_ok()); + } + + #[test] + fn extension_instance_view_maps_plugin_instance_identity() { + let instance = PluginInstanceEntry { + id: "builtin:webhook:notify:primary".to_string(), + plugin_id: "builtin:webhook".to_string(), + domain: PluginContractDomain::Notify, + subsystem: "notify_webhook".to_string(), + account_id: "primary".to_string(), + service: "webhook".to_string(), + status: "offline".to_string(), + source: PluginInstanceSource::Config, + enabled: true, + config: HashMap::from([("endpoint".to_string(), "https://example.test/webhook".to_string())]), + operational_state: Some(PluginOperationalStateContract { + install_state: PluginInstallState::Installed, + enable_state: PluginEnableState::Enabled, + runtime_state: PluginOperationalRuntimeState::Offline, + }), + diagnostic_codes: Vec::new(), + }; + + let mapped = map_extension_instance(instance.clone()); + + assert_eq!(mapped.id, instance.id); + assert_eq!(mapped.extension_id, instance.plugin_id); + assert_eq!(mapped.kind, ExtensionKind::TargetPlugin); + assert_eq!(mapped.domain, instance.domain); + assert_eq!(mapped.source, instance.source); + assert_eq!(mapped.operational_state, instance.operational_state); + } + + fn extract_block_between_markers<'a>(src: &'a str, start_marker: &str, end_marker: &str) -> &'a str { + let start = src + .find(start_marker) + .unwrap_or_else(|| panic!("Expected marker `{start_marker}` in source")); + let after_start = &src[start..]; + let end = after_start + .find(end_marker) + .unwrap_or_else(|| panic!("Expected end marker `{end_marker}` in source")); + &after_start[..end] + } +} diff --git a/rustfs/src/admin/handlers/mod.rs b/rustfs/src/admin/handlers/mod.rs index f411f8489..277a58804 100644 --- a/rustfs/src/admin/handlers/mod.rs +++ b/rustfs/src/admin/handlers/mod.rs @@ -18,6 +18,7 @@ mod audit_runtime_config; pub mod bucket_meta; pub mod config_admin; pub mod event; +pub mod extensions; pub mod group; pub mod heal; pub mod health; @@ -71,6 +72,8 @@ mod tests { let _set_config_handler = config_admin::SetConfigHandler {}; let _list_audit_targets = audit::ListAuditTargets {}; let _get_module_switches = module_switch::GetModuleSwitchesHandler {}; + let _get_extension_catalog = extensions::GetExtensionCatalogHandler {}; + let _list_extension_instances = extensions::ListExtensionInstancesHandler {}; let _get_plugin_catalog = plugins_catalog::GetPluginCatalogHandler {}; let _list_plugin_instances = plugins_instances::ListPluginInstancesHandler {}; let _get_plugin_instance = plugins_instances::GetPluginInstanceHandler {}; diff --git a/rustfs/src/admin/handlers/plugins_instances.rs b/rustfs/src/admin/handlers/plugins_instances.rs index 8103129f2..66008e067 100644 --- a/rustfs/src/admin/handlers/plugins_instances.rs +++ b/rustfs/src/admin/handlers/plugins_instances.rs @@ -286,7 +286,7 @@ struct PluginInstanceBody { } #[derive(Debug, Default, Clone, PartialEq, Eq)] -struct PluginInstanceFilters { +pub(super) struct PluginInstanceFilters { domain: Option, service: Option, status: Option, @@ -298,7 +298,7 @@ struct PluginInstanceFilters { marker: Option, } -fn extract_plugin_instance_filters(req: &S3Request) -> S3Result { +pub(super) fn extract_plugin_instance_filters(req: &S3Request) -> S3Result { let mut filters = PluginInstanceFilters::default(); if let Some(query) = req.uri.query() { @@ -420,12 +420,15 @@ fn resolve_plugin_instance_target(instance_id: &str) -> S3Result, filters: &PluginInstanceFilters) -> Vec { +pub(super) fn filter_plugin_instances( + mut instances: Vec, + filters: &PluginInstanceFilters, +) -> Vec { instances.retain(|instance| plugin_instance_matches_filters(instance, filters)); instances } -fn paginate_plugin_instances( +pub(super) fn paginate_plugin_instances( instances: Vec, filters: &PluginInstanceFilters, ) -> S3Result<(Vec, bool, Option)> { @@ -453,7 +456,7 @@ fn paginate_plugin_instances( Ok((page, truncated, next_marker)) } -fn collect_diagnostic_counts(instances: &[PluginInstanceEntry]) -> Vec { +pub(super) fn collect_diagnostic_counts(instances: &[PluginInstanceEntry]) -> Vec { let mut counts = BTreeMap::::new(); for instance in instances { for code in &instance.diagnostic_codes { @@ -628,7 +631,7 @@ async fn collect_domain_instances(context: PluginInstanceDomainContext) -> S3Res Ok(entries) } -async fn collect_all_instances() -> S3Result> { +pub(super) async fn collect_all_instances() -> S3Result> { let (mut notify_instances, audit_instances) = tokio::try_join!( collect_domain_instances(plugin_instance_domain_context(PluginContractDomain::Notify)), collect_domain_instances(plugin_instance_domain_context(PluginContractDomain::Audit)) diff --git a/rustfs/src/admin/mod.rs b/rustfs/src/admin/mod.rs index 28b5d7d99..a454bc626 100644 --- a/rustfs/src/admin/mod.rs +++ b/rustfs/src/admin/mod.rs @@ -31,8 +31,9 @@ mod console_test; mod route_registration_test; use handlers::{ - audit, bucket_meta, config_admin, heal, health, kms, module_switch, oidc, plugins_catalog, plugins_instances, pools, - profile_admin, quota, rebalance, replication, scanner, site_replication, sts, system, table_catalog, tier, tls_debug, user, + audit, bucket_meta, config_admin, extensions, heal, health, kms, module_switch, oidc, plugins_catalog, plugins_instances, + pools, profile_admin, quota, rebalance, replication, scanner, site_replication, sts, system, table_catalog, tier, tls_debug, + user, }; use router::{AdminOperation, S3Router}; use s3s::route::S3Route; @@ -69,6 +70,7 @@ fn register_admin_routes(r: &mut S3Router) -> std::io::Result<() scanner::register_scanner_route(r)?; audit::register_audit_target_route(r)?; module_switch::register_module_switch_route(r)?; + extensions::register_extension_route(r)?; plugins_catalog::register_plugin_catalog_route(r)?; plugins_instances::register_plugin_instance_route(r)?; diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index 3f3889372..2016a94e2 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -223,6 +223,8 @@ fn expected_admin_route_matrix() -> Vec { ), admin_route(Method::GET, "/v3/module-switches"), admin_route(Method::PUT, "/v3/module-switches"), + admin_route(Method::GET, "/v4/extensions/catalog"), + admin_route(Method::GET, "/v4/extensions/instances"), admin_route(Method::GET, "/v4/plugins/catalog"), admin_route(Method::GET, "/v4/plugins/instances"), admin_route_sample(Method::GET, "/v4/plugins/instances/{id}", "/v4/plugins/instances/example-id"), @@ -540,6 +542,8 @@ fn test_register_routes_cover_representative_admin_paths() { assert_route(&router, Method::GET, &admin_path("/v3/audit/target/list")); assert_route(&router, Method::GET, &admin_path("/v3/module-switches")); assert_route(&router, Method::PUT, &admin_path("/v3/module-switches")); + assert_route(&router, Method::GET, &admin_path("/v4/extensions/catalog")); + assert_route(&router, Method::GET, &admin_path("/v4/extensions/instances")); assert_route(&router, Method::GET, &admin_path("/v4/plugins/catalog")); assert_route(&router, Method::GET, &admin_path("/v4/plugins/instances")); assert_route(&router, Method::GET, &admin_path("/v4/plugins/instances/example-id"));