onetaskgraph_core/subprocess/
plugin.rs1use 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
16pub(crate) const KIND: &str = "subprocess";
18
19pub(crate) const SETTINGS_FIELD: &str = "settings";
24
25#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
34#[serde(deny_unknown_fields)]
35pub struct SubprocessConfig {
36 pub command: Program,
38 #[serde(default)]
40 pub args: Vec<String>,
41 #[serde(default)]
43 pub secrets: Vec<CredentialName>,
44 #[serde(default)]
46 pub settings: Value,
47 #[serde(default = "default_deadline_ms")]
49 pub deadline_ms: NonZeroU64,
50}
51
52fn default_deadline_ms() -> NonZeroU64 {
53 RequestDeadline::DEFAULT.milliseconds()
54}
55
56#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
63#[serde(into = "String", try_from = "String")]
64#[schemars(with = "String")]
66pub struct Program(String);
67
68impl Program {
69 #[must_use]
71 pub fn new(command: &str) -> Option<Self> {
72 (!command.trim().is_empty()).then(|| Self(command.to_owned()))
73 }
74
75 #[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#[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 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
156fn 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}