Skip to main content

onetaskgraph_core/subprocess/
plugin.rs

1//! Configuring a source that is another program.
2
3use std::collections::BTreeMap;
4use std::num::NonZeroU64;
5use std::path::Path;
6
7use onetaskgraph_plugin_api::{SecretResolver, SourceError, SourceName, SourcePlugin, TaskSource};
8use schemars::{JsonSchema, Schema, schema_for};
9use secrecy::ExposeSecret;
10use serde::{Deserialize, Serialize};
11use serde_json::Value;
12
13use super::source::{RequestDeadline, SubprocessSource};
14use crate::secrets::CredentialName;
15
16/// The name a configuration document's `plugin:` field names this kind by.
17pub(crate) const KIND: &str = "subprocess";
18
19/// The field of [`SubprocessConfig`] handed to the child as its `config:` block.
20///
21/// Named once because the configuration layer reads it too: it is the block whose
22/// supplying document the handshake reports as `document_dir`.
23pub(crate) const SETTINGS_FIELD: &str = "settings";
24
25/// How to run a plugin that speaks `docs/plugin-protocol.md`.
26///
27/// `settings` is the seam that keeps this one plugin general: it is handed to the child
28/// as its `config:` block verbatim, so what a Python source needs and what a Rust one
29/// needs are that child's business rather than a field of this schema. Which is also why
30/// the credentials a child may see are *named* here rather than inherited: §3.1 forbids a
31/// plugin reading credentials from its own environment, and forwarding the engine's whole
32/// environment would hand every plugin every secret on the host.
33#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
34#[serde(deny_unknown_fields)]
35pub struct SubprocessConfig {
36    /// The program to run.
37    pub command: Program,
38    /// Arguments to run it with.
39    #[serde(default)]
40    pub args: Vec<String>,
41    /// The environment variables whose resolved values the handshake forwards.
42    #[serde(default)]
43    pub secrets: Vec<CredentialName>,
44    /// This source's own settings, handed to the child verbatim.
45    #[serde(default)]
46    pub settings: Value,
47    /// Defaults to 30 seconds and cannot be zero.
48    #[serde(default = "default_deadline_ms")]
49    pub deadline_ms: NonZeroU64,
50}
51
52fn default_deadline_ms() -> NonZeroU64 {
53    RequestDeadline::DEFAULT.milliseconds()
54}
55
56/// The program that serves a source: a name that is not blank.
57///
58/// A newtype rather than a `String` checked on the way past, because the check has to hold
59/// wherever one of these comes from. A blank command is not a source that fails later; it
60/// is a source that was never configured, and the difference is the difference between a
61/// sentence naming the field to fill in and a spawn error about an empty path.
62#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
63#[serde(into = "String", try_from = "String")]
64// schemars does not read `serde(into)`, and this is a string on the wire.
65#[schemars(with = "String")]
66pub struct Program(String);
67
68impl Program {
69    /// `command` when it names something, and nothing otherwise.
70    #[must_use]
71    pub fn new(command: &str) -> Option<Self> {
72        (!command.trim().is_empty()).then(|| Self(command.to_owned()))
73    }
74
75    /// The program itself.
76    #[must_use]
77    pub fn as_str(&self) -> &str {
78        &self.0
79    }
80}
81
82impl TryFrom<String> for Program {
83    type Error = String;
84
85    fn try_from(value: String) -> Result<Self, Self::Error> {
86        Self::new(&value)
87            .ok_or_else(|| "`command` must name the program that serves this source".to_owned())
88    }
89}
90
91impl From<Program> for String {
92    fn from(value: Program) -> Self {
93        value.0
94    }
95}
96
97/// The factory a configuration reaches through `plugin: subprocess`.
98#[derive(Debug, Clone, Copy, Default)]
99pub struct Plugin;
100
101impl SourcePlugin for Plugin {
102    fn kind(&self) -> &'static str {
103        KIND
104    }
105
106    fn config_schema(&self) -> Schema {
107        schema_for!(SubprocessConfig)
108    }
109
110    fn build(
111        &self,
112        name: &SourceName,
113        config: &Value,
114        secrets: &dyn SecretResolver,
115    ) -> Result<Box<dyn TaskSource>, SourceError> {
116        self.build_from_document(name, config, secrets, None)
117    }
118}
119
120impl Plugin {
121    /// [`SourcePlugin::build`], telling the child which document supplied its settings.
122    ///
123    /// `document_dir` is [`SourceConfig::document_dir`](crate::SourceConfig::document_dir):
124    /// the engine reaches this rather than the trait method because the trait hands a
125    /// plugin values and no origins, and this plugin alone passes an origin on — as
126    /// `document_dir` in the handshake (`docs/plugin-protocol.md` §3).
127    ///
128    /// # Errors
129    ///
130    /// What [`SourcePlugin::build`] returns.
131    pub fn build_from_document(
132        &self,
133        name: &SourceName,
134        config: &Value,
135        secrets: &dyn SecretResolver,
136        document_dir: Option<&Path>,
137    ) -> Result<Box<dyn TaskSource>, SourceError> {
138        let config: SubprocessConfig =
139            serde_json::from_value(config.clone()).map_err(|error| SourceError::Config {
140                message: format!("source {name}: {error}"),
141            })?;
142        let forwarded = resolve_named(name, &config.secrets, secrets)?;
143        SubprocessSource::connect_from_document(
144            config.command.as_str(),
145            &config.args,
146            name,
147            &config.settings,
148            forwarded,
149            RequestDeadline::from_millis(config.deadline_ms),
150            document_dir,
151        )
152        .map(|source| Box::new(source) as Box<dyn TaskSource>)
153    }
154}
155
156/// The values of exactly the variables this configuration names, and nothing else.
157///
158/// A named variable nothing defines is refused here rather than forwarded as absent: the
159/// plugin asked for it, so a run that spawned anyway would fail later inside the child
160/// with a message about a credential, and the thing the user has to fix is on this side.
161fn resolve_named(
162    name: &SourceName,
163    named: &[CredentialName],
164    secrets: &dyn SecretResolver,
165) -> Result<BTreeMap<String, String>, SourceError> {
166    let mut forwarded = BTreeMap::new();
167    for variable in named {
168        let value = secrets
169            .get(variable.as_str())
170            .ok_or_else(|| SourceError::Auth {
171                message: format!(
172                    "source {name}: nothing defines {variable}, which this source's \
173                     `secrets` names; export it, or add it to the credentials file"
174                ),
175            })?;
176        forwarded.insert(
177            variable.as_str().to_owned(),
178            value.expose_secret().to_owned(),
179        );
180    }
181    Ok(forwarded)
182}