Skip to main content

r2x_python/
plugin_invoker.rs

1//! Plugin invocation and execution
2
3use crate::errors::BridgeError;
4use r2x_logger as logger;
5use r2x_manifest::runtime::{build_runtime_bindings, PluginRole, RuntimeBindings};
6use r2x_manifest::types::Plugin;
7use std::fs;
8use std::io::ErrorKind;
9use std::path::{Component, Path, PathBuf};
10use std::time::Duration;
11
12/// A directory-backed plugin artifact with one JSON entrypoint.
13///
14/// System artifacts can add sidecar files next to the entrypoint. Consumers
15/// must therefore retain and pass the complete bundle directory.
16#[derive(Clone, Debug, PartialEq, Eq)]
17pub struct ArtifactBundle {
18    root: PathBuf,
19    entrypoint: PathBuf,
20}
21
22impl ArtifactBundle {
23    /// Create a bundle rooted at `root` with a relative JSON entrypoint.
24    pub fn new(
25        root: impl Into<PathBuf>,
26        entrypoint: impl Into<PathBuf>,
27    ) -> Result<Self, BridgeError> {
28        let entrypoint = entrypoint.into();
29        validate_relative_entrypoint(&entrypoint)?;
30        Ok(Self {
31            root: root.into(),
32            entrypoint,
33        })
34    }
35
36    /// Root directory containing the entrypoint and any sidecars.
37    pub fn root(&self) -> &Path {
38        &self.root
39    }
40
41    /// Relative path of the JSON entrypoint within [`Self::root`].
42    pub fn relative_entrypoint(&self) -> &Path {
43        &self.entrypoint
44    }
45
46    /// Absolute or relative filesystem path to the JSON entrypoint.
47    pub fn entrypoint_path(&self) -> PathBuf {
48        self.root.join(&self.entrypoint)
49    }
50}
51
52fn validate_relative_entrypoint(entrypoint: &Path) -> Result<(), BridgeError> {
53    if entrypoint.as_os_str().is_empty() {
54        return Err(BridgeError::InvalidArtifact(
55            "entrypoint cannot be empty".to_string(),
56        ));
57    }
58
59    let mut has_file_name = false;
60    for component in entrypoint.components() {
61        match component {
62            Component::Normal(_) => has_file_name = true,
63            Component::CurDir => {
64                return Err(BridgeError::InvalidArtifact(format!(
65                    "entrypoint cannot contain '.' components: {}",
66                    entrypoint.display()
67                )));
68            }
69            Component::ParentDir => {
70                return Err(BridgeError::InvalidArtifact(format!(
71                    "entrypoint cannot traverse its bundle root: {}",
72                    entrypoint.display()
73                )));
74            }
75            Component::RootDir | Component::Prefix(_) => {
76                return Err(BridgeError::InvalidArtifact(format!(
77                    "entrypoint must be relative to its bundle root: {}",
78                    entrypoint.display()
79                )));
80            }
81        }
82    }
83    if !has_file_name {
84        return Err(BridgeError::InvalidArtifact(
85            "entrypoint must name a file within its bundle".to_string(),
86        ));
87    }
88    Ok(())
89}
90
91/// Type of data produced by an artifact-mode plugin invocation.
92#[derive(Clone, Copy, Debug, PartialEq, Eq)]
93pub enum ArtifactOutputKind {
94    /// A `System` persisted through `System.to_json(path)`.
95    System,
96    /// Generic JSON persisted through Python's JSON backend.
97    Json,
98    /// No replacement artifact was emitted.
99    Empty,
100}
101
102/// Result of artifact-mode plugin invocation.
103#[derive(Debug)]
104pub struct PluginArtifactInvocationResult {
105    /// The materialized output type, if any.
106    pub output_kind: ArtifactOutputKind,
107    /// Optional per-phase timings for diagnostics.
108    pub timings: Option<PluginInvocationTimings>,
109}
110
111/// Timings for a plugin invocation phase
112#[derive(Debug)]
113pub struct PluginInvocationTimings {
114    pub python_invocation: Duration,
115    pub serialization: Duration,
116}
117
118/// Direct input supplied to a plugin invocation.
119#[derive(Clone, Copy, Debug)]
120pub enum PluginInput<'a> {
121    /// JSON payload received from standard input.
122    Json(&'a str),
123    /// JSON entrypoint on disk, loaded relative to its sidecar bundle.
124    File(&'a Path),
125}
126
127/// Materialized result of running a plugin through the Python bridge.
128#[derive(Debug)]
129pub enum PluginInvocationOutput {
130    /// JSON text that should be emitted or written by the caller.
131    Json(String),
132    /// A System was written directly to the requested output path.
133    Persisted,
134    /// The plugin intentionally produced no stream output.
135    Empty,
136}
137
138/// Result of running a plugin through the Python bridge.
139#[derive(Debug)]
140pub struct PluginInvocationResult {
141    /// Materialized plugin output.
142    pub output: PluginInvocationOutput,
143    /// Optional per-phase timings for diagnostics.
144    pub timings: Option<PluginInvocationTimings>,
145}
146
147impl crate::python_bridge::Bridge {
148    /// Invoke a plugin through the direct CLI interface.
149    ///
150    /// Direct invocations treat an unused stream input as an error, redirect
151    /// plugin writes away from stdout, and can persist System results directly
152    /// to a durable JSON entrypoint.
153    pub fn invoke_plugin_direct(
154        &self,
155        target: &str,
156        config_json: &str,
157        input: Option<PluginInput<'_>>,
158        output_path: Option<&Path>,
159        plugin_metadata: Option<&Plugin>,
160    ) -> Result<PluginInvocationResult, BridgeError> {
161        let runtime_bindings = plugin_metadata.map(build_runtime_bindings);
162
163        if runtime_bindings
164            .as_ref()
165            .is_some_and(|bindings| bindings.role == PluginRole::Upgrader)
166        {
167            if input.is_some() {
168                return Err(BridgeError::Stream(
169                    "upgrader plugins do not accept System input".to_string(),
170                ));
171            }
172            return self.invoke_upgrader_plugin(
173                target,
174                config_json,
175                runtime_bindings.as_ref(),
176                plugin_metadata,
177                true,
178            );
179        }
180
181        self.invoke_plugin_regular_direct(
182            target,
183            config_json,
184            input,
185            output_path,
186            runtime_bindings.as_ref(),
187        )
188    }
189
190    pub fn invoke_plugin_with_bindings(
191        &self,
192        target: &str,
193        config_json: &str,
194        stdin_json: Option<&str>,
195        runtime_bindings: Option<&RuntimeBindings>,
196    ) -> Result<PluginInvocationResult, BridgeError> {
197        if let Some(bindings) = runtime_bindings {
198            if bindings.role == PluginRole::Upgrader {
199                logger::debug("Routing to upgrader plugin handler (runtime bindings)");
200                return self.invoke_upgrader_plugin(
201                    target,
202                    config_json,
203                    Some(bindings),
204                    None,
205                    false,
206                );
207            }
208        }
209
210        self.invoke_plugin_regular(target, config_json, stdin_json, runtime_bindings)
211    }
212
213    /// Save a System artifact as an infrasys ZIP archive.
214    pub fn save_system_artifact_as_zip(
215        &self,
216        input: &ArtifactBundle,
217        output: &Path,
218    ) -> Result<(), BridgeError> {
219        Self::save_system_artifact_as_zip_native(input, output)
220    }
221
222    /// Artifact-mode counterpart of [`Self::invoke_plugin_with_bindings`].
223    pub fn invoke_plugin_with_artifact_bindings(
224        &self,
225        target: &str,
226        config_json: &str,
227        input: Option<&ArtifactBundle>,
228        output: &ArtifactBundle,
229        runtime_bindings: Option<&RuntimeBindings>,
230    ) -> Result<PluginArtifactInvocationResult, BridgeError> {
231        if runtime_bindings.is_some_and(|bindings| bindings.role == PluginRole::Upgrader) {
232            return Err(BridgeError::UnsupportedArtifactMode(
233                "upgrader plugins are not yet supported because registered SYSTEM steps still serialize payloads through Rust".to_string(),
234            ));
235        }
236
237        validate_output_bundle(input, output)?;
238
239        self.invoke_plugin_regular_with_artifacts(
240            target,
241            config_json,
242            input,
243            output,
244            runtime_bindings,
245        )
246    }
247}
248
249fn validate_output_bundle(
250    input: Option<&ArtifactBundle>,
251    output: &ArtifactBundle,
252) -> Result<(), BridgeError> {
253    if input.is_some_and(|input| input.root() == output.root()) {
254        return Err(BridgeError::InvalidArtifact(
255            "input and output bundles must use different roots".to_string(),
256        ));
257    }
258
259    let metadata = match fs::symlink_metadata(output.root()) {
260        Ok(metadata) => metadata,
261        Err(error) if error.kind() == ErrorKind::NotFound => return Ok(()),
262        Err(error) => return Err(error.into()),
263    };
264    if metadata.file_type().is_symlink() || !metadata.is_dir() {
265        return Err(BridgeError::InvalidArtifact(format!(
266            "output bundle root must be a directory: {}",
267            output.root().display()
268        )));
269    }
270    if fs::read_dir(output.root())?.next().transpose()?.is_some() {
271        return Err(BridgeError::InvalidArtifact(format!(
272            "output bundle root must be empty: {}",
273            output.root().display()
274        )));
275    }
276
277    Ok(())
278}
279
280#[cfg(test)]
281mod tests {
282    use crate::plugin_invoker::*;
283    use crate::python_bridge::Bridge;
284    use r2x_manifest::runtime::{PluginRole, RuntimeBindings};
285    use r2x_manifest::types::PluginType;
286    use std::error::Error;
287    use tempfile::tempdir;
288
289    #[test]
290    fn plugin_invocation_result_basics() {
291        let result = PluginInvocationResult {
292            output: PluginInvocationOutput::Empty,
293            timings: None,
294        };
295        assert!(matches!(result.output, PluginInvocationOutput::Empty));
296    }
297
298    #[test]
299    fn artifact_bundle_rejects_absolute_and_traversing_entrypoints() {
300        let absolute = ArtifactBundle::new("bundle", "/tmp/system.json");
301        assert!(absolute.is_err());
302
303        let traversal = ArtifactBundle::new("bundle", "../system.json");
304        assert!(traversal.is_err());
305
306        let bundle = ArtifactBundle::new("bundle", "nested/system.json");
307        assert!(bundle.is_ok());
308
309        let directory = ArtifactBundle::new("bundle", ".");
310        assert!(directory.is_err());
311
312        let current_directory = ArtifactBundle::new("bundle", "./system.json");
313        assert!(current_directory.is_err());
314    }
315
316    #[test]
317    fn artifact_mode_rejects_upgraders_until_their_payload_path_is_native(
318    ) -> Result<(), BridgeError> {
319        let bridge = Bridge::for_tests();
320        let output = ArtifactBundle::new("bundle", "system.json")?;
321        let bindings = RuntimeBindings {
322            entry_module: "plugin".to_string(),
323            entry_name: "Upgrader".to_string(),
324            plugin_type: PluginType::Class,
325            role: PluginRole::Upgrader,
326            call_method: Some("run".to_string()),
327            config: None,
328            parameters: Vec::new(),
329            requires_store: false,
330        };
331
332        let error = bridge.invoke_plugin_with_artifact_bindings(
333            "plugin:Upgrader",
334            "{}",
335            None,
336            &output,
337            Some(&bindings),
338        );
339        assert!(matches!(
340            error,
341            Err(BridgeError::UnsupportedArtifactMode(_))
342        ));
343        Ok(())
344    }
345
346    #[test]
347    fn artifact_mode_rejects_nonempty_output_bundles() -> Result<(), Box<dyn Error>> {
348        let temp = tempdir()?;
349        let output_root = temp.path().join("output");
350        std::fs::create_dir_all(&output_root)?;
351        std::fs::write(output_root.join("stale.h5"), "stale")?;
352        let output = ArtifactBundle::new(&output_root, "system.json")?;
353
354        let result = Bridge::for_tests().invoke_plugin_with_artifact_bindings(
355            "missing:plugin",
356            "{}",
357            None,
358            &output,
359            None,
360        );
361
362        assert!(matches!(result, Err(BridgeError::InvalidArtifact(_))));
363        assert!(output_root.join("stale.h5").exists());
364        Ok(())
365    }
366}