nemo_relay/plugin/dynamic/
registry.rs1use std::collections::BTreeMap;
5
6use super::{
7 DynamicPluginAttestationMode, DynamicPluginCheckState, DynamicPluginFailure, DynamicPluginId,
8 DynamicPluginManifest, DynamicPluginMetadata, DynamicPluginRecord, DynamicPluginRuntimeStatus,
9 DynamicPluginStartupClass, DynamicPluginValidationStatus, bump_generation,
10 stamp_creation_metadata,
11};
12use crate::plugin::{PluginError, Result};
13
14#[derive(Debug, Default)]
16pub struct DynamicPluginRegistry {
17 records: BTreeMap<DynamicPluginId, DynamicPluginRecord>,
18}
19
20impl DynamicPluginRegistry {
21 pub fn new() -> Self {
23 Self::default()
24 }
25
26 pub fn from_records(records: Vec<DynamicPluginRecord>) -> Result<Self> {
28 let mut registry = Self::new();
29 for mut record in records {
30 normalize_record_shape(&mut record);
31 validate_record_shape(&record)?;
32 let plugin_id = record.metadata.id.clone();
33 if registry.records.contains_key(&plugin_id) {
34 return Err(PluginError::Conflict(format!(
35 "dynamic plugin '{plugin_id}' is duplicated in persisted registry state"
36 )));
37 }
38 registry.records.insert(plugin_id, record);
39 }
40 Ok(registry)
41 }
42
43 pub fn get(&self, plugin_id: &str) -> Option<&DynamicPluginRecord> {
45 self.records.get(plugin_id)
46 }
47
48 pub fn list(&self, include_tombstoned: bool) -> Vec<&DynamicPluginRecord> {
50 self.records
51 .values()
52 .filter(|record| include_tombstoned || !record.is_tombstoned())
53 .collect()
54 }
55
56 pub fn cloned_records(&self, include_tombstoned: bool) -> Vec<DynamicPluginRecord> {
58 self.list(include_tombstoned).into_iter().cloned().collect()
59 }
60
61 pub fn add(&mut self, mut record: DynamicPluginRecord) -> Result<&DynamicPluginRecord> {
67 normalize_record_shape(&mut record);
68 validate_record_shape(&record)?;
69
70 let plugin_id = record.metadata.id.clone();
71 record.spec.present = true;
72
73 if let Some(existing) = self.records.get(&plugin_id) {
74 if !existing.is_tombstoned() {
75 return Err(PluginError::Conflict(format!(
76 "dynamic plugin '{plugin_id}' is already registered"
77 )));
78 }
79
80 inherit_tombstoned_lineage(&mut record.metadata, &existing.metadata);
81 }
82
83 stamp_creation_metadata(&mut record.metadata);
84
85 self.records.insert(plugin_id.clone(), record);
86 Ok(self
87 .records
88 .get(&plugin_id)
89 .expect("dynamic plugin record must exist immediately after insert"))
90 }
91
92 pub fn add_manifest(
94 &mut self,
95 manifest: DynamicPluginManifest,
96 manifest_ref: Option<String>,
97 ) -> Result<&DynamicPluginRecord> {
98 let record = manifest.into_record(manifest_ref)?;
99 self.add(record)
100 }
101
102 pub fn enable(&mut self, plugin_id: &str) -> Result<bool> {
104 let record = self.lookup_mut(plugin_id)?;
105 ensure_live_record(record, plugin_id)?;
106 if record.spec.enabled {
107 return Ok(false);
108 }
109 record.spec.enabled = true;
110 bump_generation(record);
111 Ok(true)
112 }
113
114 pub fn disable(&mut self, plugin_id: &str) -> Result<bool> {
116 let record = self.lookup_mut(plugin_id)?;
117 ensure_live_record(record, plugin_id)?;
118 if !record.spec.enabled {
119 return Ok(false);
120 }
121 record.spec.enabled = false;
122 bump_generation(record);
123 Ok(true)
124 }
125
126 pub fn remove(&mut self, plugin_id: &str) -> Result<bool> {
128 let record = self.lookup_mut(plugin_id)?;
129 if record.is_tombstoned() {
130 return Ok(false);
131 }
132 record.spec.present = false;
133 record.spec.enabled = false;
134 bump_generation(record);
135 Ok(true)
136 }
137
138 pub fn update_validation_status(
140 &mut self,
141 plugin_id: &str,
142 mut validation: DynamicPluginValidationStatus,
143 ) -> Result<()> {
144 validation.checked_at = Some(super::current_timestamp());
145 let record = self.lookup_mut(plugin_id)?;
146 record.status.validation = validation;
147 Ok(())
148 }
149
150 pub fn update_runtime_status(
152 &mut self,
153 plugin_id: &str,
154 mut runtime: DynamicPluginRuntimeStatus,
155 ) -> Result<()> {
156 runtime.updated_at = Some(super::current_timestamp());
157 let record = self.lookup_mut(plugin_id)?;
158 record.status.runtime = runtime;
159 Ok(())
160 }
161
162 pub fn update_last_error(
164 &mut self,
165 plugin_id: &str,
166 last_error: Option<DynamicPluginFailure>,
167 ) -> Result<()> {
168 let record = self.lookup_mut(plugin_id)?;
169 record.status.last_error = last_error;
170 Ok(())
171 }
172
173 pub fn update_environment(
175 &mut self,
176 plugin_id: &str,
177 environment_ref: Option<String>,
178 environment: DynamicPluginCheckState,
179 ) -> Result<()> {
180 let record = self.lookup_mut(plugin_id)?;
181 record.source.environment_ref = environment_ref;
182 record.status.validation.environment = environment;
183 record.status.validation.checked_at = Some(super::current_timestamp());
184 Ok(())
185 }
186
187 pub fn update_policy_status(
189 &mut self,
190 plugin_id: &str,
191 policy_satisfied: DynamicPluginCheckState,
192 startup_class: DynamicPluginStartupClass,
193 attestation_mode: DynamicPluginAttestationMode,
194 last_error: Option<DynamicPluginFailure>,
195 ) -> Result<()> {
196 let record = self.lookup_mut(plugin_id)?;
197 record.status.validation.policy_satisfied = policy_satisfied;
198 record.status.startup_class = Some(startup_class);
199 record.status.attestation_mode = Some(attestation_mode);
200 record.status.last_error = last_error;
201 Ok(())
202 }
203
204 fn lookup_mut(&mut self, plugin_id: &str) -> Result<&mut DynamicPluginRecord> {
205 self.records.get_mut(plugin_id).ok_or_else(|| {
206 PluginError::NotFound(format!("dynamic plugin '{plugin_id}' is not registered"))
207 })
208 }
209}
210
211fn ensure_live_record(record: &DynamicPluginRecord, plugin_id: &str) -> Result<()> {
212 if record.is_tombstoned() {
213 return Err(PluginError::Conflict(format!(
214 "dynamic plugin '{plugin_id}' has been removed"
215 )));
216 }
217 Ok(())
218}
219
220fn inherit_tombstoned_lineage(
221 metadata: &mut DynamicPluginMetadata,
222 existing: &DynamicPluginMetadata,
223) {
224 let next_generation = existing.generation.saturating_add(1);
225 if metadata.created_at.is_none() {
226 metadata.created_at = existing.created_at.clone();
227 }
228 metadata.generation = next_generation;
229}
230
231fn normalize_record_shape(record: &mut DynamicPluginRecord) {
232 record.metadata.id = record.metadata.id.trim().to_owned();
233 match &mut record.compatibility {
234 super::DynamicPluginCompatibility::RustDynamic(compatibility) => {
235 compatibility.relay = compatibility.relay.trim().to_owned();
236 compatibility.native_api = compatibility.native_api.trim().to_owned();
237 }
238 super::DynamicPluginCompatibility::Worker(compatibility) => {
239 compatibility.relay = compatibility.relay.trim().to_owned();
240 compatibility.worker_protocol = compatibility.worker_protocol.trim().to_owned();
241 }
242 }
243
244 match &mut record.load {
245 super::DynamicPluginLoadContract::Worker(load) => {
246 load.entrypoint = load.entrypoint.trim().to_owned();
247 }
248 super::DynamicPluginLoadContract::RustDynamic(load) => {
249 load.library = load.library.trim().to_owned();
250 load.symbol = load.symbol.trim().to_owned();
251 }
252 }
253}
254
255fn validate_record_shape(record: &DynamicPluginRecord) -> Result<()> {
256 if record.metadata.id.trim().is_empty() {
257 return Err(PluginError::InvalidConfig(
258 "dynamic plugin id must not be empty".into(),
259 ));
260 }
261
262 match record.metadata.kind {
263 super::DynamicPluginKind::RustDynamic => validate_rust_dynamic_record(record),
264 super::DynamicPluginKind::Worker => validate_worker_record(record),
265 }
266}
267
268fn validate_rust_dynamic_record(record: &DynamicPluginRecord) -> Result<()> {
269 let super::DynamicPluginCompatibility::RustDynamic(compatibility) = &record.compatibility
270 else {
271 return Err(PluginError::InvalidConfig(
272 "dynamic rust_dynamic record has invalid compatibility shape".into(),
273 ));
274 };
275 validate_relay_compatibility(&compatibility.relay)?;
276 if compatibility.native_api.trim().is_empty() {
277 return Err(PluginError::InvalidConfig(
278 "dynamic rust_dynamic record has invalid compatibility shape".into(),
279 ));
280 }
281 let valid_load = matches!(
282 &record.load,
283 super::DynamicPluginLoadContract::RustDynamic(load)
284 if !load.library.trim().is_empty() && !load.symbol.trim().is_empty()
285 );
286 if valid_load {
287 Ok(())
288 } else {
289 Err(PluginError::InvalidConfig(
290 "dynamic rust_dynamic record has invalid load shape".into(),
291 ))
292 }
293}
294
295fn validate_worker_record(record: &DynamicPluginRecord) -> Result<()> {
296 let super::DynamicPluginCompatibility::Worker(compatibility) = &record.compatibility else {
297 return Err(PluginError::InvalidConfig(
298 "dynamic worker record has invalid compatibility shape".into(),
299 ));
300 };
301 validate_relay_compatibility(&compatibility.relay)?;
302 if compatibility.worker_protocol.trim().is_empty() {
303 return Err(PluginError::InvalidConfig(
304 "dynamic worker record has invalid compatibility shape".into(),
305 ));
306 }
307 let valid_load = matches!(
308 &record.load,
309 super::DynamicPluginLoadContract::Worker(load) if !load.entrypoint.trim().is_empty()
310 );
311 if valid_load {
312 Ok(())
313 } else {
314 Err(PluginError::InvalidConfig(
315 "dynamic worker record has invalid load shape".into(),
316 ))
317 }
318}
319
320fn validate_relay_compatibility(relay: &str) -> Result<()> {
321 if relay.trim().is_empty() {
322 Err(PluginError::InvalidConfig(
323 "dynamic plugin record must declare compat.relay".into(),
324 ))
325 } else {
326 Ok(())
327 }
328}