1pub mod docs {
54 #[doc = include_str!("../docs/plugin-contract.md")]
56 pub mod plugin_contract {}
57
58 #[doc = include_str!("../docs/build-your-own-plugin.md")]
60 pub mod build_your_own_plugin {}
61
62 #[doc = include_str!("../docs/echo-plugin.md")]
64 pub mod echo_plugin {}
65}
66
67use std::collections::HashMap;
68use std::time::Duration;
69
70use flare_grpc_proto::capability::capability_service_client::CapabilityServiceClient;
71use flare_grpc_proto::capability::{
72 DeregisterPluginEndpointRequest, RegisterPluginEndpointRequest,
73};
74use tonic::transport::Channel;
75
76#[derive(Debug, thiserror::Error)]
77pub enum HostError {
78 #[error("连接 capability 失败:{0}")]
79 Connect(String),
80 #[error("capability 拒绝了注册:{0}")]
81 Rejected(String),
82 #[error("gRPC 调用失败:{0}")]
83 Rpc(String),
84 #[error("声明不完整:{0}")]
85 Invalid(&'static str),
86}
87
88#[derive(Debug, Clone)]
94pub struct PluginDeclaration {
95 pub tenant_id: String,
96 pub plugin_id: String,
97 pub capability_id: String,
103 pub grpc_authority: String,
105 pub plugin_version: String,
106 pub api_version: String,
107 pub manifest_sha256: String,
109 pub declared_operations: Vec<String>,
111 pub labels: HashMap<String, String>,
113 pub seat_model: SeatModel,
118}
119
120#[derive(Debug, Clone, Copy, PartialEq, Eq)]
127pub enum SeatModel {
128 Tenant,
129 PerUser,
130}
131
132impl SeatModel {
133 fn as_wire(self) -> &'static str {
134 match self {
135 Self::Tenant => "tenant",
136 Self::PerUser => "per_user",
137 }
138 }
139}
140
141impl PluginDeclaration {
142 pub fn validate(&self) -> Result<(), HostError> {
147 if self.plugin_id.trim().is_empty() {
148 return Err(HostError::Invalid("plugin_id 不能为空"));
149 }
150 if self.capability_id.trim().is_empty() {
151 return Err(HostError::Invalid("capability_id 不能为空"));
152 }
153 if self.grpc_authority.trim().is_empty() {
154 return Err(HostError::Invalid("grpc_authority 不能为空"));
155 }
156 if self.declared_operations.is_empty() {
157 return Err(HostError::Invalid(
158 "declared_operations 为空会让插件退化为 unverified,声明边界无法强制",
159 ));
160 }
161 if !self.declared_operations.contains(&self.capability_id) {
162 return Err(HostError::Invalid(
163 "capability_id 必须出现在 declared_operations 里",
164 ));
165 }
166 Ok(())
167 }
168}
169
170pub struct PluginHost {
172 client: CapabilityServiceClient<Channel>,
173}
174
175impl PluginHost {
176 pub async fn connect(endpoint: impl Into<String>) -> Result<Self, HostError> {
178 let endpoint = endpoint.into();
179 let channel = Channel::from_shared(endpoint.clone())
180 .map_err(|e| HostError::Connect(format!("{endpoint}: {e}")))?
181 .connect_timeout(Duration::from_secs(5))
182 .connect()
183 .await
184 .map_err(|e| HostError::Connect(format!("{endpoint}: {e}")))?;
185 Ok(Self {
186 client: CapabilityServiceClient::new(channel),
187 })
188 }
189
190 pub async fn announce(&mut self, declaration: &PluginDeclaration) -> Result<(), HostError> {
194 declaration.validate()?;
195
196 let response = self
197 .client
198 .register_plugin_endpoint(RegisterPluginEndpointRequest {
199 tenant_id: declaration.tenant_id.clone(),
200 plugin_id: declaration.plugin_id.clone(),
201 capability_id: declaration.capability_id.clone(),
202 grpc_authority: declaration.grpc_authority.clone(),
203 labels: declaration.labels.clone(),
204 request_id: String::new(),
205 plugin_version: declaration.plugin_version.clone(),
206 api_version: declaration.api_version.clone(),
207 manifest_sha256: declaration.manifest_sha256.clone(),
208 declared_operations: declaration.declared_operations.clone(),
209 seat_model: declaration.seat_model.as_wire().to_string(),
210 })
211 .await
212 .map_err(|e| HostError::Rpc(e.to_string()))?
213 .into_inner();
214
215 if !response.accepted {
216 return Err(HostError::Rejected(response.message));
217 }
218 tracing::info!(
219 plugin_id = %declaration.plugin_id,
220 capability_id = %declaration.capability_id,
221 operations = declaration.declared_operations.len(),
222 "plugin announced to capability"
223 );
224 Ok(())
225 }
226
227 pub async fn withdraw(&mut self, tenant_id: &str, plugin_id: &str) -> Result<(), HostError> {
234 self.client
235 .deregister_plugin_endpoint(DeregisterPluginEndpointRequest {
236 tenant_id: tenant_id.to_string(),
237 plugin_id: plugin_id.to_string(),
238 request_id: String::new(),
239 })
240 .await
241 .map_err(|e| HostError::Rpc(e.to_string()))?;
242 tracing::info!(plugin_id, "plugin withdrawn from capability");
243 Ok(())
244 }
245}
246
247#[cfg(test)]
248mod tests {
249 use super::*;
250
251 fn declaration() -> PluginDeclaration {
252 PluginDeclaration {
253 tenant_id: "0".into(),
254 plugin_id: "p1".into(),
255 capability_id: "vendorx.do".into(),
256 grpc_authority: "127.0.0.1:1".into(),
257 plugin_version: "1.0.0".into(),
258 api_version: "1".into(),
259 manifest_sha256: "abc".into(),
260 declared_operations: vec!["vendorx.do".into()],
261 labels: HashMap::new(),
262 seat_model: SeatModel::Tenant,
263 }
264 }
265
266 #[test]
267 fn complete_declaration_is_valid() {
268 declaration().validate().expect("完整声明应当通过");
269 }
270
271 #[test]
273 fn empty_declared_operations_is_rejected_locally() {
274 let mut d = declaration();
275 d.declared_operations.clear();
276 assert!(matches!(d.validate(), Err(HostError::Invalid(_))));
277 }
278
279 #[test]
281 fn capability_id_must_be_declared() {
282 let mut d = declaration();
283 d.declared_operations = vec!["vendorx.other".into()];
284 assert!(matches!(d.validate(), Err(HostError::Invalid(_))));
285 }
286}