Skip to main content

nemo_relay/plugin/
dynamic.rs

1// SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Dynamic plugin control-plane and registry model.
5//!
6//! This module owns the durable control-plane record shape for dynamic plugins.
7//! Authored manifest parsing/validation and in-memory registry mutation logic
8//! live in dedicated submodules so those responsibilities do not accumulate in
9//! one file as the feature grows.
10
11use chrono::Utc;
12use semver::{Version, VersionReq};
13use serde::{Deserialize, Serialize};
14use strum::{Display, IntoStaticStr};
15
16use crate::plugin::{
17    PluginDeregistrationOutcome, PluginError, deregister_plugin_registration_checked,
18};
19
20/// Canonical identifier for one dynamic plugin record.
21pub type DynamicPluginId = String;
22
23/// Canonical filename for authored Relay plugin manifests.
24pub const DYNAMIC_PLUGIN_MANIFEST_FILENAME: &str = "relay-plugin.toml";
25
26mod host;
27mod manifest;
28mod native;
29mod registry;
30#[cfg(feature = "worker-grpc")]
31mod worker;
32
33pub use host::*;
34pub use manifest::*;
35pub use native::*;
36pub use registry::*;
37#[cfg(feature = "worker-grpc")]
38pub use worker::*;
39
40#[derive(Debug)]
41pub(crate) struct DynamicPluginTeardownOutcome {
42    pub(crate) errors: Vec<String>,
43    pub(crate) safe_to_unload: bool,
44}
45
46impl DynamicPluginTeardownOutcome {
47    pub(crate) fn success() -> Self {
48        Self {
49            errors: Vec::new(),
50            safe_to_unload: true,
51        }
52    }
53
54    pub(crate) fn record_error(&mut self, error: impl Into<String>, safe_to_unload: bool) {
55        self.errors.push(error.into());
56        self.safe_to_unload &= safe_to_unload;
57    }
58
59    pub(crate) fn merge(&mut self, other: Self) {
60        self.errors.extend(other.errors);
61        self.safe_to_unload &= other.safe_to_unload;
62    }
63}
64
65pub(super) fn deregister_tracked_registrations_checked(
66    registrations: &mut Vec<(String, u64)>,
67    plugin_type: &str,
68) -> DynamicPluginTeardownOutcome {
69    let mut outcome = DynamicPluginTeardownOutcome::success();
70    for (plugin_kind, registration_id) in std::mem::take(registrations).into_iter().rev() {
71        match deregister_plugin_registration_checked(&plugin_kind, registration_id) {
72            Ok(PluginDeregistrationOutcome::Removed) => {}
73            Ok(PluginDeregistrationOutcome::Missing) => outcome.record_error(
74                format!(
75                    "{plugin_type} plugin kind '{plugin_kind}' was not registered during teardown"
76                ),
77                true,
78            ),
79            Ok(PluginDeregistrationOutcome::Replaced) => outcome.record_error(
80                format!(
81                    "{plugin_type} plugin kind '{plugin_kind}' was replaced during teardown and was left registered"
82                ),
83                true,
84            ),
85            Err(error) => outcome.record_error(
86                format!(
87                    "failed to deregister {plugin_type} plugin kind '{plugin_kind}': {error}"
88                ),
89                false,
90            ),
91        }
92    }
93    outcome
94}
95
96pub(super) fn validate_annotated_request_consumer_compatibility(
97    relay: &str,
98    plugin_kind: &str,
99) -> crate::plugin::Result<()> {
100    let requirement = VersionReq::parse(relay).map_err(|error| {
101        PluginError::InvalidConfig(format!("invalid compat.relay version requirement: {error}"))
102    })?;
103    if requirement.matches(&Version::new(0, 5, u64::MAX)) {
104        return Err(PluginError::InvalidConfig(format!(
105            "dynamic plugin '{plugin_kind}' registers an LLM request intercept and must declare compat.relay = \">=0.6,<1.0\" or another range that excludes Relay 0.5"
106        )));
107    }
108    Ok(())
109}
110
111/// Plugin execution lane.
112#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash, Display)]
113#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
114#[serde(rename_all = "snake_case")]
115#[strum(serialize_all = "snake_case")]
116pub enum DynamicPluginKind {
117    /// Trusted in-process native plugin.
118    RustDynamic,
119    /// Isolated worker-based plugin runtime.
120    Worker,
121}
122
123/// Managed runtime identity for worker-based plugins.
124#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash, Display)]
125#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
126#[serde(rename_all = "snake_case")]
127#[strum(serialize_all = "snake_case")]
128pub enum WorkerRuntime {
129    /// Python worker runtime.
130    Python,
131    /// Rust worker executable runtime.
132    Rust,
133    /// Generic executable worker runtime.
134    Command,
135}
136
137/// Relay-enforced capability declared by a dynamic plugin.
138#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash, Display)]
139#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
140#[serde(rename_all = "snake_case")]
141#[strum(serialize_all = "snake_case")]
142pub enum DynamicPluginCapability {
143    /// Trusted in-process native extension capability.
144    PluginNative,
145    /// Isolated worker-based extension capability.
146    PluginWorker,
147    /// Typed configuration schema contribution capability.
148    ConfigSchema,
149}
150
151/// Host policy startup classification for a plugin.
152#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash, Display)]
153#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
154#[serde(rename_all = "snake_case")]
155#[strum(serialize_all = "snake_case")]
156pub enum DynamicPluginStartupClass {
157    /// Failure is tolerated and the host may start in degraded mode.
158    Optional,
159    /// Failure is startup-fatal under current host policy.
160    Required,
161}
162
163/// Host attestation policy mode.
164#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash, Display)]
165#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
166#[serde(rename_all = "snake_case")]
167#[strum(serialize_all = "snake_case")]
168pub enum DynamicPluginAttestationMode {
169    /// Integrity verification only.
170    IntegrityOnly,
171    /// Verify signatures when present but do not require them.
172    SignatureIfPresent,
173    /// Require trusted signature verification.
174    SignatureRequired,
175}
176
177/// High-level verification state for one validation axis.
178#[derive(
179    Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq, Hash, IntoStaticStr,
180)]
181#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
182#[serde(rename_all = "snake_case")]
183#[strum(serialize_all = "snake_case")]
184pub enum DynamicPluginCheckState {
185    /// No verification result is currently known.
186    #[default]
187    Unknown,
188    /// Verification passed.
189    Valid,
190    /// Verification failed.
191    Invalid,
192}
193
194/// Observed runtime state for a dynamic plugin.
195#[derive(
196    Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq, Hash, IntoStaticStr,
197)]
198#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
199#[serde(rename_all = "snake_case")]
200#[strum(serialize_all = "snake_case")]
201pub enum DynamicPluginRuntimeState {
202    /// Not currently active.
203    #[default]
204    Stopped,
205    /// Activation is in progress.
206    Starting,
207    /// Currently active.
208    Running,
209    /// Activation failed or the active runtime crashed.
210    Failed,
211}
212
213/// Recent failure phase for diagnostics and operator UX.
214#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Hash, Display)]
215#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
216#[serde(rename_all = "snake_case")]
217#[strum(serialize_all = "snake_case")]
218pub enum DynamicPluginFailurePhase {
219    /// Failure occurred during validation.
220    Validation,
221    /// Failure occurred during activation or reconciliation.
222    Activation,
223    /// Failure occurred after activation while running.
224    Runtime,
225    /// Failure occurred because policy no longer permits realization.
226    Policy,
227}
228
229/// Stable metadata for one durable plugin record.
230#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
231#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
232pub struct DynamicPluginMetadata {
233    /// Canonical plugin identifier.
234    pub id: DynamicPluginId,
235    /// Optional human-friendly display label.
236    #[serde(default, skip_serializing_if = "Option::is_none")]
237    pub name: Option<String>,
238    /// Optional plugin version mirrored from packaging metadata when desired.
239    #[serde(default, skip_serializing_if = "Option::is_none")]
240    pub version: Option<String>,
241    /// Execution lane used by the plugin.
242    pub kind: DynamicPluginKind,
243    /// Monotonic desired-state generation.
244    #[serde(default)]
245    pub generation: u64,
246    /// Creation timestamp in RFC 3339 form.
247    #[serde(default, skip_serializing_if = "Option::is_none")]
248    pub created_at: Option<String>,
249    /// Last durable record update time in RFC 3339 form.
250    #[serde(default, skip_serializing_if = "Option::is_none")]
251    pub updated_at: Option<String>,
252}
253
254/// Source and resolved artifact facts for a plugin.
255#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
256#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
257pub struct DynamicPluginSource {
258    /// Canonical manifest location or reference.
259    #[serde(default, skip_serializing_if = "Option::is_none")]
260    pub manifest_ref: Option<String>,
261    /// Resolved runtime artifact location.
262    #[serde(default, skip_serializing_if = "Option::is_none")]
263    pub artifact_ref: Option<String>,
264    /// Resolved environment location for worker-based plugins.
265    #[serde(default, skip_serializing_if = "Option::is_none")]
266    pub environment_ref: Option<String>,
267    /// Pinned artifact digest.
268    #[serde(default, skip_serializing_if = "Option::is_none")]
269    pub artifact_digest: Option<String>,
270}
271
272/// Desired-state fields owned by user-facing operations.
273#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
274#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
275pub struct DynamicPluginSpec {
276    /// Whether the plugin should be present in desired state.
277    #[serde(default = "default_present")]
278    pub present: bool,
279    /// Whether the plugin should be enabled in desired state.
280    #[serde(default)]
281    pub enabled: bool,
282    /// Optional config reference controlled by higher-level config surfaces.
283    #[serde(default, skip_serializing_if = "Option::is_none")]
284    pub config_ref: Option<String>,
285}
286
287pub(crate) fn default_present() -> bool {
288    true
289}
290
291impl Default for DynamicPluginSpec {
292    fn default() -> Self {
293        Self {
294            present: true,
295            enabled: false,
296            config_ref: None,
297        }
298    }
299}
300
301/// Lane-specific compatibility declarations and resolved compatibility facts.
302#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
303#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
304#[serde(rename_all = "snake_case")]
305pub enum DynamicPluginCompatibility {
306    /// Native shared-library compatibility contract.
307    RustDynamic(DynamicPluginRustCompatibility),
308    /// Worker runtime compatibility contract.
309    Worker(DynamicPluginWorkerCompatibility),
310}
311
312/// Compatibility contract for worker plugins.
313#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
314#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
315pub struct DynamicPluginWorkerCompatibility {
316    /// Compatible NeMo Relay version or version range.
317    pub relay: String,
318    /// Worker protocol version for `worker`.
319    pub worker_protocol: String,
320}
321
322/// Compatibility contract for native shared libraries.
323#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
324#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
325pub struct DynamicPluginRustCompatibility {
326    /// Compatible NeMo Relay version or version range.
327    pub relay: String,
328    /// Native host API/ABI contract version for `rust_dynamic`.
329    pub native_api: String,
330}
331
332/// Runtime entry contract for the resolved plugin.
333#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
334#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
335#[serde(rename_all = "snake_case")]
336pub enum DynamicPluginLoadContract {
337    /// Worker-based plugin registration target.
338    Worker(DynamicPluginWorkerLoadContract),
339    /// Native shared-library registration target.
340    RustDynamic(DynamicPluginRustLoadContract),
341}
342
343/// Lane-specific load contract for worker plugins.
344#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
345#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
346pub struct DynamicPluginWorkerLoadContract {
347    /// Managed worker runtime identity.
348    pub runtime: WorkerRuntime,
349    /// Worker entrypoint or registration target.
350    pub entrypoint: String,
351}
352
353/// Lane-specific load contract for native shared libraries.
354#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
355#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
356pub struct DynamicPluginRustLoadContract {
357    /// Native dynamic library path.
358    pub library: String,
359    /// Native exported registration symbol.
360    pub symbol: String,
361}
362
363/// One structured recent failure summary.
364#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
365#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
366pub struct DynamicPluginFailure {
367    /// Failure phase.
368    pub phase: DynamicPluginFailurePhase,
369    /// Machine-readable failure code.
370    pub code: String,
371    /// Human-readable summary.
372    pub message: String,
373}
374
375/// Decomposed validation results for one plugin record.
376#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
377#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
378pub struct DynamicPluginValidationStatus {
379    /// Manifest schema/result state.
380    #[serde(default)]
381    pub manifest: DynamicPluginCheckState,
382    /// Relay/native/worker compatibility state.
383    #[serde(default)]
384    pub compatibility: DynamicPluginCheckState,
385    /// Artifact integrity state.
386    #[serde(default)]
387    pub integrity: DynamicPluginCheckState,
388    /// Environment/runtime readiness state.
389    #[serde(default)]
390    pub environment: DynamicPluginCheckState,
391    /// Signature/authenticity state.
392    #[serde(default)]
393    pub authenticity: DynamicPluginCheckState,
394    /// Whether the current host policy is satisfied.
395    #[serde(default)]
396    pub policy_satisfied: DynamicPluginCheckState,
397    /// Most recent validation time in RFC 3339 form.
398    #[serde(default, skip_serializing_if = "Option::is_none")]
399    pub checked_at: Option<String>,
400    /// Concise operator-facing validation summary.
401    #[serde(default, skip_serializing_if = "Option::is_none")]
402    pub message: Option<String>,
403}
404
405impl Default for DynamicPluginValidationStatus {
406    fn default() -> Self {
407        Self {
408            manifest: DynamicPluginCheckState::Unknown,
409            compatibility: DynamicPluginCheckState::Unknown,
410            integrity: DynamicPluginCheckState::Unknown,
411            environment: DynamicPluginCheckState::Unknown,
412            authenticity: DynamicPluginCheckState::Unknown,
413            policy_satisfied: DynamicPluginCheckState::Unknown,
414            checked_at: None,
415            message: None,
416        }
417    }
418}
419
420/// Observed runtime state for one plugin record.
421#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
422#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
423pub struct DynamicPluginRuntimeStatus {
424    /// Current observed runtime state.
425    #[serde(default)]
426    pub state: DynamicPluginRuntimeState,
427    /// Desired-state generation this runtime status reflects.
428    #[serde(default)]
429    pub observed_generation: u64,
430    /// Most recent successful start/activation time.
431    #[serde(default, skip_serializing_if = "Option::is_none")]
432    pub started_at: Option<String>,
433    /// Most recent runtime-status refresh time.
434    #[serde(default, skip_serializing_if = "Option::is_none")]
435    pub updated_at: Option<String>,
436    /// Concise operator-facing runtime summary.
437    #[serde(default, skip_serializing_if = "Option::is_none")]
438    pub message: Option<String>,
439}
440
441impl Default for DynamicPluginRuntimeStatus {
442    fn default() -> Self {
443        Self {
444            state: DynamicPluginRuntimeState::Stopped,
445            observed_generation: 0,
446            started_at: None,
447            updated_at: None,
448            message: None,
449        }
450    }
451}
452
453/// Durable observed state for a plugin record.
454#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
455#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
456pub struct DynamicPluginStatus {
457    /// Validation and policy status.
458    #[serde(default)]
459    pub validation: DynamicPluginValidationStatus,
460    /// Runtime state observed by the control plane.
461    #[serde(default)]
462    pub runtime: DynamicPluginRuntimeStatus,
463    /// Host policy startup classification.
464    #[serde(default, skip_serializing_if = "Option::is_none")]
465    pub startup_class: Option<DynamicPluginStartupClass>,
466    /// Effective attestation mode for this plugin under host policy.
467    #[serde(default, skip_serializing_if = "Option::is_none")]
468    pub attestation_mode: Option<DynamicPluginAttestationMode>,
469    /// Most recent meaningful failure summary.
470    #[serde(default, skip_serializing_if = "Option::is_none")]
471    pub last_error: Option<DynamicPluginFailure>,
472}
473
474/// Durable control-plane record for a dynamic plugin.
475#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
476#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
477pub struct DynamicPluginRecord {
478    /// Stable plugin metadata.
479    pub metadata: DynamicPluginMetadata,
480    /// Source and artifact facts.
481    #[serde(default)]
482    pub source: DynamicPluginSource,
483    /// Desired state.
484    #[serde(default)]
485    pub spec: DynamicPluginSpec,
486    /// Compatibility declarations and resolved compatibility facts.
487    pub compatibility: DynamicPluginCompatibility,
488    /// Resolved runtime entry contract.
489    pub load: DynamicPluginLoadContract,
490    /// Observed state.
491    #[serde(default)]
492    pub status: DynamicPluginStatus,
493}
494
495impl DynamicPluginRecord {
496    /// Returns `true` when the runtime has observed the current desired-state generation.
497    pub fn is_reconciled(&self) -> bool {
498        self.status.runtime.observed_generation == self.metadata.generation
499    }
500
501    /// Returns `true` when the record is tombstoned.
502    pub fn is_tombstoned(&self) -> bool {
503        !self.spec.present
504    }
505}
506
507pub(crate) fn current_timestamp() -> String {
508    Utc::now().to_rfc3339()
509}
510
511pub(crate) fn stamp_creation_metadata(metadata: &mut DynamicPluginMetadata) {
512    if metadata.created_at.is_none() {
513        metadata.created_at = Some(current_timestamp());
514    }
515    if metadata.updated_at.is_none() {
516        metadata.updated_at = metadata.created_at.clone();
517    }
518}
519
520pub(crate) fn touch_metadata(metadata: &mut DynamicPluginMetadata) {
521    metadata.updated_at = Some(current_timestamp());
522}
523
524pub(crate) fn bump_generation(record: &mut DynamicPluginRecord) {
525    record.metadata.generation = record.metadata.generation.saturating_add(1);
526    touch_metadata(&mut record.metadata);
527}
528
529#[cfg(test)]
530#[path = "../../tests/unit/plugin_dynamic_tests.rs"]
531mod tests;