Skip to main content

nemo_relay/plugin/dynamic/
registry.rs

1// SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4use 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/// In-memory dynamic plugin registry used by the control plane.
15#[derive(Debug, Default)]
16pub struct DynamicPluginRegistry {
17    records: BTreeMap<DynamicPluginId, DynamicPluginRecord>,
18}
19
20impl DynamicPluginRegistry {
21    /// Creates an empty dynamic plugin registry.
22    pub fn new() -> Self {
23        Self::default()
24    }
25
26    /// Reconstructs a registry from previously persisted durable records.
27    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    /// Returns the registered record for `plugin_id`, if present.
44    pub fn get(&self, plugin_id: &str) -> Option<&DynamicPluginRecord> {
45        self.records.get(plugin_id)
46    }
47
48    /// Lists records, hiding tombstones unless requested.
49    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    /// Clones records for serialization or higher-level projection.
57    pub fn cloned_records(&self, include_tombstoned: bool) -> Vec<DynamicPluginRecord> {
58        self.list(include_tombstoned).into_iter().cloned().collect()
59    }
60
61    /// Adds a new dynamic plugin record.
62    ///
63    /// This is a trusted internal control-plane API. Callers that start from an
64    /// authored `relay-plugin.toml` manifest should prefer [`Self::add_manifest`]
65    /// so the manifest contract is enforced before record creation.
66    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    /// Validates an authored manifest and registers the resulting dynamic plugin record.
93    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    /// Marks the plugin enabled in desired state.
103    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    /// Marks the plugin disabled in desired state.
115    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    /// Tombstones the plugin record and disables desired runtime realization.
127    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    /// Replaces the current validation status without mutating desired state.
139    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    /// Replaces the current runtime status without mutating desired state.
151    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    /// Records the most recent dynamic-plugin failure summary.
163    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    /// Replaces the resolved worker environment and its validation state.
174    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    /// Replaces the current host-policy outcome without mutating desired state.
188    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}