1use std::collections::HashMap;
50use std::time::Duration;
51
52use flare_grpc_proto::capability::capability_service_client::CapabilityServiceClient;
53use flare_grpc_proto::capability::{
54 DeregisterPluginEndpointRequest, RegisterPluginEndpointRequest,
55};
56use tonic::transport::Channel;
57
58#[derive(Debug, thiserror::Error)]
59pub enum HostError {
60 #[error("连接 capability 失败:{0}")]
61 Connect(String),
62 #[error("capability 拒绝了注册:{0}")]
63 Rejected(String),
64 #[error("gRPC 调用失败:{0}")]
65 Rpc(String),
66 #[error("声明不完整:{0}")]
67 Invalid(&'static str),
68}
69
70#[derive(Debug, Clone)]
76pub struct PluginDeclaration {
77 pub tenant_id: String,
78 pub plugin_id: String,
79 pub capability_id: String,
85 pub grpc_authority: String,
87 pub plugin_version: String,
88 pub api_version: String,
89 pub manifest_sha256: String,
91 pub declared_operations: Vec<String>,
93 pub labels: HashMap<String, String>,
95 pub seat_model: SeatModel,
100}
101
102#[derive(Debug, Clone, Copy, PartialEq, Eq)]
109pub enum SeatModel {
110 Tenant,
111 PerUser,
112}
113
114impl SeatModel {
115 fn as_wire(self) -> &'static str {
116 match self {
117 Self::Tenant => "tenant",
118 Self::PerUser => "per_user",
119 }
120 }
121}
122
123impl PluginDeclaration {
124 pub fn validate(&self) -> Result<(), HostError> {
129 if self.plugin_id.trim().is_empty() {
130 return Err(HostError::Invalid("plugin_id 不能为空"));
131 }
132 if self.capability_id.trim().is_empty() {
133 return Err(HostError::Invalid("capability_id 不能为空"));
134 }
135 if self.grpc_authority.trim().is_empty() {
136 return Err(HostError::Invalid("grpc_authority 不能为空"));
137 }
138 if self.declared_operations.is_empty() {
139 return Err(HostError::Invalid(
140 "declared_operations 为空会让插件退化为 unverified,声明边界无法强制",
141 ));
142 }
143 if !self.declared_operations.contains(&self.capability_id) {
144 return Err(HostError::Invalid(
145 "capability_id 必须出现在 declared_operations 里",
146 ));
147 }
148 Ok(())
149 }
150}
151
152pub struct PluginHost {
154 client: CapabilityServiceClient<Channel>,
155}
156
157impl PluginHost {
158 pub async fn connect(endpoint: impl Into<String>) -> Result<Self, HostError> {
160 let endpoint = endpoint.into();
161 let channel = Channel::from_shared(endpoint.clone())
162 .map_err(|e| HostError::Connect(format!("{endpoint}: {e}")))?
163 .connect_timeout(Duration::from_secs(5))
164 .connect()
165 .await
166 .map_err(|e| HostError::Connect(format!("{endpoint}: {e}")))?;
167 Ok(Self {
168 client: CapabilityServiceClient::new(channel),
169 })
170 }
171
172 pub async fn announce(&mut self, declaration: &PluginDeclaration) -> Result<(), HostError> {
176 declaration.validate()?;
177
178 let response = self
179 .client
180 .register_plugin_endpoint(RegisterPluginEndpointRequest {
181 tenant_id: declaration.tenant_id.clone(),
182 plugin_id: declaration.plugin_id.clone(),
183 capability_id: declaration.capability_id.clone(),
184 grpc_authority: declaration.grpc_authority.clone(),
185 labels: declaration.labels.clone(),
186 request_id: String::new(),
187 plugin_version: declaration.plugin_version.clone(),
188 api_version: declaration.api_version.clone(),
189 manifest_sha256: declaration.manifest_sha256.clone(),
190 declared_operations: declaration.declared_operations.clone(),
191 seat_model: declaration.seat_model.as_wire().to_string(),
192 })
193 .await
194 .map_err(|e| HostError::Rpc(e.to_string()))?
195 .into_inner();
196
197 if !response.accepted {
198 return Err(HostError::Rejected(response.message));
199 }
200 tracing::info!(
201 plugin_id = %declaration.plugin_id,
202 capability_id = %declaration.capability_id,
203 operations = declaration.declared_operations.len(),
204 "plugin announced to capability"
205 );
206 Ok(())
207 }
208
209 pub async fn withdraw(&mut self, tenant_id: &str, plugin_id: &str) -> Result<(), HostError> {
216 self.client
217 .deregister_plugin_endpoint(DeregisterPluginEndpointRequest {
218 tenant_id: tenant_id.to_string(),
219 plugin_id: plugin_id.to_string(),
220 request_id: String::new(),
221 })
222 .await
223 .map_err(|e| HostError::Rpc(e.to_string()))?;
224 tracing::info!(plugin_id, "plugin withdrawn from capability");
225 Ok(())
226 }
227}
228
229#[cfg(test)]
230mod tests {
231 use super::*;
232
233 fn declaration() -> PluginDeclaration {
234 PluginDeclaration {
235 tenant_id: "0".into(),
236 plugin_id: "p1".into(),
237 capability_id: "vendorx.do".into(),
238 grpc_authority: "127.0.0.1:1".into(),
239 plugin_version: "1.0.0".into(),
240 api_version: "1".into(),
241 manifest_sha256: "abc".into(),
242 declared_operations: vec!["vendorx.do".into()],
243 labels: HashMap::new(),
244 seat_model: SeatModel::Tenant,
245 }
246 }
247
248 #[test]
249 fn complete_declaration_is_valid() {
250 declaration().validate().expect("完整声明应当通过");
251 }
252
253 #[test]
255 fn empty_declared_operations_is_rejected_locally() {
256 let mut d = declaration();
257 d.declared_operations.clear();
258 assert!(matches!(d.validate(), Err(HostError::Invalid(_))));
259 }
260
261 #[test]
263 fn capability_id_must_be_declared() {
264 let mut d = declaration();
265 d.declared_operations = vec!["vendorx.other".into()];
266 assert!(matches!(d.validate(), Err(HostError::Invalid(_))));
267 }
268}