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/// Options that control one plugin invocation.
128#[derive(Clone, Copy, Debug, Default)]
129pub struct PluginInvocationOptions {
130    /// Enter Python's post-mortem debugger when an uncaught plugin exception occurs.
131    pub pdb: bool,
132}
133
134/// Materialized result of running a plugin through the Python bridge.
135#[derive(Debug)]
136pub enum PluginInvocationOutput {
137    /// JSON text that should be emitted or written by the caller.
138    Json(String),
139    /// A System was written directly to the requested output path.
140    Persisted,
141    /// The plugin intentionally produced no stream output.
142    Empty,
143}
144
145/// Result of running a plugin through the Python bridge.
146#[derive(Debug)]
147pub struct PluginInvocationResult {
148    /// Materialized plugin output.
149    pub output: PluginInvocationOutput,
150    /// Optional per-phase timings for diagnostics.
151    pub timings: Option<PluginInvocationTimings>,
152}
153
154impl crate::python_bridge::Bridge {
155    /// Invoke a plugin through the direct CLI interface.
156    ///
157    /// Direct invocations treat an unused stream input as an error, redirect
158    /// plugin writes away from stdout, and can persist System results directly
159    /// to a durable JSON entrypoint.
160    pub fn invoke_plugin_direct(
161        &self,
162        target: &str,
163        config_json: &str,
164        input: Option<PluginInput<'_>>,
165        output_path: Option<&Path>,
166        plugin_metadata: Option<&Plugin>,
167        options: PluginInvocationOptions,
168    ) -> Result<PluginInvocationResult, BridgeError> {
169        let runtime_bindings = plugin_metadata.map(build_runtime_bindings);
170
171        if runtime_bindings
172            .as_ref()
173            .is_some_and(|bindings| bindings.role == PluginRole::Upgrader)
174        {
175            if input.is_some() {
176                return Err(BridgeError::Stream(
177                    "upgrader plugins do not accept System input".to_string(),
178                ));
179            }
180            return self.invoke_upgrader_plugin(
181                target,
182                config_json,
183                runtime_bindings.as_ref(),
184                plugin_metadata,
185                true,
186                options,
187            );
188        }
189
190        self.invoke_plugin_regular_direct(
191            target,
192            config_json,
193            input,
194            output_path,
195            runtime_bindings.as_ref(),
196            options,
197        )
198    }
199
200    pub fn invoke_plugin_with_bindings(
201        &self,
202        target: &str,
203        config_json: &str,
204        stdin_json: Option<&str>,
205        runtime_bindings: Option<&RuntimeBindings>,
206        options: PluginInvocationOptions,
207    ) -> Result<PluginInvocationResult, BridgeError> {
208        if let Some(bindings) = runtime_bindings {
209            if bindings.role == PluginRole::Upgrader {
210                logger::debug("Routing to upgrader plugin handler (runtime bindings)");
211                return self.invoke_upgrader_plugin(
212                    target,
213                    config_json,
214                    Some(bindings),
215                    None,
216                    false,
217                    options,
218                );
219            }
220        }
221
222        self.invoke_plugin_regular(target, config_json, stdin_json, runtime_bindings, options)
223    }
224
225    /// Save a System artifact as an infrasys ZIP archive.
226    pub fn save_system_artifact_as_zip(
227        &self,
228        input: &ArtifactBundle,
229        output: &Path,
230    ) -> Result<(), BridgeError> {
231        Self::save_system_artifact_as_zip_native(input, output)
232    }
233
234    /// Artifact-mode counterpart of [`Self::invoke_plugin_with_bindings`].
235    pub fn invoke_plugin_with_artifact_bindings(
236        &self,
237        target: &str,
238        config_json: &str,
239        input: Option<&ArtifactBundle>,
240        output: &ArtifactBundle,
241        runtime_bindings: Option<&RuntimeBindings>,
242        options: PluginInvocationOptions,
243    ) -> Result<PluginArtifactInvocationResult, BridgeError> {
244        if runtime_bindings.is_some_and(|bindings| bindings.role == PluginRole::Upgrader) {
245            return Err(BridgeError::UnsupportedArtifactMode(
246                "upgrader plugins are not yet supported because registered SYSTEM steps still serialize payloads through Rust".to_string(),
247            ));
248        }
249
250        validate_output_bundle(input, output)?;
251
252        self.invoke_plugin_regular_with_artifacts(
253            target,
254            config_json,
255            input,
256            output,
257            runtime_bindings,
258            options,
259        )
260    }
261}
262
263fn validate_output_bundle(
264    input: Option<&ArtifactBundle>,
265    output: &ArtifactBundle,
266) -> Result<(), BridgeError> {
267    if input.is_some_and(|input| input.root() == output.root()) {
268        return Err(BridgeError::InvalidArtifact(
269            "input and output bundles must use different roots".to_string(),
270        ));
271    }
272
273    let metadata = match fs::symlink_metadata(output.root()) {
274        Ok(metadata) => metadata,
275        Err(error) if error.kind() == ErrorKind::NotFound => return Ok(()),
276        Err(error) => return Err(error.into()),
277    };
278    if metadata.file_type().is_symlink() || !metadata.is_dir() {
279        return Err(BridgeError::InvalidArtifact(format!(
280            "output bundle root must be a directory: {}",
281            output.root().display()
282        )));
283    }
284    if fs::read_dir(output.root())?.next().transpose()?.is_some() {
285        return Err(BridgeError::InvalidArtifact(format!(
286            "output bundle root must be empty: {}",
287            output.root().display()
288        )));
289    }
290
291    Ok(())
292}
293
294#[cfg(test)]
295mod tests {
296    use crate::plugin_invoker::*;
297    use crate::python_bridge::Bridge;
298    use r2x_manifest::runtime::{PluginRole, RuntimeBindings};
299    use r2x_manifest::types::PluginType;
300    use std::error::Error;
301    use tempfile::tempdir;
302
303    #[test]
304    fn plugin_invocation_result_basics() {
305        let result = PluginInvocationResult {
306            output: PluginInvocationOutput::Empty,
307            timings: None,
308        };
309        assert!(matches!(result.output, PluginInvocationOutput::Empty));
310    }
311
312    #[test]
313    fn artifact_bundle_rejects_absolute_and_traversing_entrypoints() {
314        let absolute = ArtifactBundle::new("bundle", "/tmp/system.json");
315        assert!(absolute.is_err());
316
317        let traversal = ArtifactBundle::new("bundle", "../system.json");
318        assert!(traversal.is_err());
319
320        let bundle = ArtifactBundle::new("bundle", "nested/system.json");
321        assert!(bundle.is_ok());
322
323        let directory = ArtifactBundle::new("bundle", ".");
324        assert!(directory.is_err());
325
326        let current_directory = ArtifactBundle::new("bundle", "./system.json");
327        assert!(current_directory.is_err());
328    }
329
330    #[test]
331    fn artifact_mode_rejects_upgraders_until_their_payload_path_is_native(
332    ) -> Result<(), BridgeError> {
333        let bridge = Bridge::for_tests();
334        let output = ArtifactBundle::new("bundle", "system.json")?;
335        let bindings = RuntimeBindings {
336            entry_module: "plugin".to_string(),
337            entry_name: "Upgrader".to_string(),
338            plugin_type: PluginType::Class,
339            role: PluginRole::Upgrader,
340            call_method: Some("run".to_string()),
341            config: None,
342            parameters: Vec::new(),
343            requires_store: false,
344        };
345
346        let error = bridge.invoke_plugin_with_artifact_bindings(
347            "plugin:Upgrader",
348            "{}",
349            None,
350            &output,
351            Some(&bindings),
352            PluginInvocationOptions::default(),
353        );
354        assert!(matches!(
355            error,
356            Err(BridgeError::UnsupportedArtifactMode(_))
357        ));
358        Ok(())
359    }
360
361    #[test]
362    fn artifact_mode_rejects_nonempty_output_bundles() -> Result<(), Box<dyn Error>> {
363        let temp = tempdir()?;
364        let output_root = temp.path().join("output");
365        std::fs::create_dir_all(&output_root)?;
366        std::fs::write(output_root.join("stale.h5"), "stale")?;
367        let output = ArtifactBundle::new(&output_root, "system.json")?;
368
369        let result = Bridge::for_tests().invoke_plugin_with_artifact_bindings(
370            "missing:plugin",
371            "{}",
372            None,
373            &output,
374            None,
375            PluginInvocationOptions::default(),
376        );
377
378        assert!(matches!(result, Err(BridgeError::InvalidArtifact(_))));
379        assert!(output_root.join("stale.h5").exists());
380        Ok(())
381    }
382}