From 1e7c376586079f4fbbdc7297ab2e1cf0708889ff 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:27:55 +0800 Subject: [PATCH] feat(targets): add extension schema adapter (#3385) --- Cargo.lock | 1 + crates/targets/Cargo.toml | 1 + crates/targets/src/catalog/extension.rs | 179 ++++++++++++++++++++++++ crates/targets/src/catalog/mod.rs | 1 + crates/targets/src/lib.rs | 4 + 5 files changed, 186 insertions(+) create mode 100644 crates/targets/src/catalog/extension.rs diff --git a/Cargo.lock b/Cargo.lock index 1178f62df..c5b9eafc6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -10031,6 +10031,7 @@ dependencies = [ "rumqttc-next", "rustfs-config", "rustfs-ecstore", + "rustfs-extension-schema", "rustfs-kafka-async", "rustfs-s3-types", "rustfs-tls-runtime", diff --git a/crates/targets/Cargo.toml b/crates/targets/Cargo.toml index dafe71bd7..191100667 100644 --- a/crates/targets/Cargo.toml +++ b/crates/targets/Cargo.toml @@ -14,6 +14,7 @@ documentation = "https://docs.rs/rustfs-targets/latest/rustfs_targets/" [dependencies] rustfs-config = { workspace = true, features = ["notify", "constants", "audit", "server-config-model"] } rustfs-ecstore = { workspace = true } +rustfs-extension-schema = { workspace = true } rustfs-tls-runtime = { workspace = true } rustfs-s3-types = { workspace = true } async-trait = { workspace = true } diff --git a/crates/targets/src/catalog/extension.rs b/crates/targets/src/catalog/extension.rs new file mode 100644 index 000000000..ab1ffd61b --- /dev/null +++ b/crates/targets/src/catalog/extension.rs @@ -0,0 +1,179 @@ +// 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 std::collections::BTreeMap; + +use crate::catalog::builtin::{builtin_audit_target_admin_descriptors, builtin_notify_target_admin_descriptors}; +use crate::domain::TargetDomain; +use crate::manifest::{ + TargetPluginEntrypointKind, TargetPluginMarketplaceManifest, TargetPluginPackaging, builtin_target_marketplace_manifest, +}; +use rustfs_extension_schema::{ + EXTENSION_SCHEMA_VERSION, ExtensionCapabilityRef, ExtensionKind, ExtensionRuntimeBoundary, ExtensionRuntimeContract, + ExtensionSchema, +}; + +pub const TARGET_AUDIT_CAPABILITY: &str = "target.audit.v1"; +pub const TARGET_NOTIFY_CAPABILITY: &str = "target.notify.v1"; + +pub fn target_marketplace_extension_schema(manifest: &TargetPluginMarketplaceManifest) -> ExtensionSchema { + let boundary = target_runtime_boundary(manifest.entrypoint_kind); + + ExtensionSchema { + schema_version: EXTENSION_SCHEMA_VERSION.to_string(), + extension_id: manifest.plugin_id.to_string(), + display_name: manifest.display_name.to_string(), + provider: manifest.provider.to_string(), + version: manifest.version.to_string(), + kind: ExtensionKind::TargetPlugin, + runtime: ExtensionRuntimeContract { + api_version: manifest.api_compatibility_version.to_string(), + boundary, + }, + capabilities: manifest + .supported_domains + .iter() + .copied() + .map(target_domain_capability) + .collect(), + disabled_by_default: target_disabled_by_default(manifest.packaging, boundary), + } +} + +pub fn builtin_target_extension_schemas() -> Vec { + let mut target_types = BTreeMap::new(); + + for descriptor in builtin_audit_target_admin_descriptors() + .into_iter() + .chain(builtin_notify_target_admin_descriptors()) + { + target_types + .entry(descriptor.manifest().target_type) + .or_insert_with(|| builtin_target_marketplace_manifest(descriptor.manifest().target_type)); + } + + target_types.values().map(target_marketplace_extension_schema).collect() +} + +pub const fn target_runtime_boundary(entrypoint_kind: TargetPluginEntrypointKind) -> ExtensionRuntimeBoundary { + match entrypoint_kind { + TargetPluginEntrypointKind::Builtin => ExtensionRuntimeBoundary::Builtin, + TargetPluginEntrypointKind::Sidecar => ExtensionRuntimeBoundary::Sidecar, + TargetPluginEntrypointKind::Wasm => ExtensionRuntimeBoundary::Wasm, + } +} + +const fn target_disabled_by_default(packaging: TargetPluginPackaging, boundary: ExtensionRuntimeBoundary) -> bool { + matches!(packaging, TargetPluginPackaging::External) || boundary.requires_disabled_by_default() +} + +fn target_domain_capability(domain: TargetDomain) -> ExtensionCapabilityRef { + match domain { + TargetDomain::Audit => ExtensionCapabilityRef::new(TARGET_AUDIT_CAPABILITY), + TargetDomain::Notify => ExtensionCapabilityRef::new(TARGET_NOTIFY_CAPABILITY), + } +} + +#[cfg(test)] +mod tests { + use super::{ + TARGET_AUDIT_CAPABILITY, TARGET_NOTIFY_CAPABILITY, builtin_target_extension_schemas, target_marketplace_extension_schema, + }; + use crate::catalog::example_external_webhook_plugin; + use crate::domain::TargetDomain; + use crate::manifest::{ + TargetPluginExternalRuntimeContract, TargetPluginPackaging, TargetPluginRuntimeTransport, + builtin_target_marketplace_manifest, + }; + use rustfs_extension_schema::{ + EXTENSION_SCHEMA_VERSION, ExtensionKind, ExtensionRuntimeBoundary, validate_extension_schemas, + }; + + #[test] + fn builtin_target_marketplace_manifest_maps_to_builtin_extension_schema() { + let schema = target_marketplace_extension_schema(&builtin_target_marketplace_manifest("webhook")); + + assert_eq!(schema.schema_version, EXTENSION_SCHEMA_VERSION); + assert_eq!(schema.extension_id, "builtin:webhook"); + assert_eq!(schema.display_name, "Webhook"); + assert_eq!(schema.kind, ExtensionKind::TargetPlugin); + assert_eq!(schema.runtime.boundary, ExtensionRuntimeBoundary::Builtin); + assert!(!schema.disabled_by_default); + assert_eq!(schema.runtime.api_version, "rustfs.target-plugin.v1"); + assert_eq!(schema.capabilities.len(), 2); + assert_eq!(schema.capabilities[0].as_str(), TARGET_AUDIT_CAPABILITY); + assert_eq!(schema.capabilities[1].as_str(), TARGET_NOTIFY_CAPABILITY); + } + + #[test] + fn external_sidecar_target_manifest_maps_to_disabled_sidecar_extension_schema() { + let example = example_external_webhook_plugin(); + let schema = target_marketplace_extension_schema(&example.manifest); + + assert_eq!(schema.extension_id, "external:webhook-sidecar"); + assert_eq!(schema.kind, ExtensionKind::TargetPlugin); + assert_eq!(schema.runtime.boundary, ExtensionRuntimeBoundary::Sidecar); + assert!(schema.disabled_by_default); + assert_eq!(schema.capabilities.len(), 1); + assert_eq!(schema.capabilities[0].as_str(), TARGET_NOTIFY_CAPABILITY); + assert!(validate_extension_schemas(&[schema]).is_ok()); + } + + #[test] + fn external_packaging_stays_disabled_even_with_builtin_entrypoint_metadata() { + let mut manifest = builtin_target_marketplace_manifest("webhook"); + manifest.packaging = TargetPluginPackaging::External; + manifest.runtime_contract = TargetPluginExternalRuntimeContract { + protocol_version: "rustfs.target-runtime.v1", + transport: TargetPluginRuntimeTransport::InProcess, + }; + manifest.supported_domains = &[TargetDomain::Notify]; + + let schema = target_marketplace_extension_schema(&manifest); + + assert_eq!(schema.runtime.boundary, ExtensionRuntimeBoundary::Builtin); + assert!(schema.disabled_by_default); + } + + #[test] + fn builtin_target_extension_catalog_is_unique_and_schema_valid() { + let schemas = builtin_target_extension_schemas(); + let mut extension_ids = schemas.iter().map(|schema| schema.extension_id.as_str()).collect::>(); + extension_ids.sort_unstable(); + + assert_eq!(schemas.len(), 9); + assert_eq!( + extension_ids, + vec![ + "builtin:amqp", + "builtin:kafka", + "builtin:mqtt", + "builtin:mysql", + "builtin:nats", + "builtin:postgres", + "builtin:pulsar", + "builtin:redis", + "builtin:webhook", + ] + ); + assert!(schemas.iter().all(|schema| schema.kind == ExtensionKind::TargetPlugin)); + assert!( + schemas + .iter() + .all(|schema| schema.runtime.boundary == ExtensionRuntimeBoundary::Builtin) + ); + assert!(schemas.iter().all(|schema| !schema.disabled_by_default)); + assert!(validate_extension_schemas(&schemas).is_ok()); + } +} diff --git a/crates/targets/src/catalog/mod.rs b/crates/targets/src/catalog/mod.rs index a6e049f5c..8c899f175 100644 --- a/crates/targets/src/catalog/mod.rs +++ b/crates/targets/src/catalog/mod.rs @@ -13,6 +13,7 @@ // limitations under the License. pub mod builtin; +pub mod extension; use crate::control_plane::external_target_plugin_installation; use crate::domain::TargetDomain; diff --git a/crates/targets/src/lib.rs b/crates/targets/src/lib.rs index 0959fecb6..cc9078ac9 100644 --- a/crates/targets/src/lib.rs +++ b/crates/targets/src/lib.rs @@ -27,6 +27,10 @@ pub mod store; pub mod sys; pub mod target; +pub use catalog::extension::{ + TARGET_AUDIT_CAPABILITY, TARGET_NOTIFY_CAPABILITY, builtin_target_extension_schemas, target_marketplace_extension_schema, + target_runtime_boundary, +}; pub use check::{ check_amqp_broker_available, check_kafka_broker_available, check_mqtt_broker_available, check_mqtt_broker_available_with_tls, check_mysql_server_available, check_nats_server_available, check_postgres_server_available, check_pulsar_broker_available,