Skip to main content

lenso_secrets_command_plugin/
lib.rs

1//! Bounded external-command Secrets Provider Plugin.
2
3use std::{
4    cell::RefCell,
5    collections::{BTreeMap, BTreeSet},
6    env, fmt, fs,
7    io::Read,
8    path::PathBuf,
9    process::{Child, Command, Stdio},
10    rc::Rc,
11    thread,
12    time::{Duration, Instant},
13};
14
15use futures::channel::oneshot;
16use lenso::prelude::*;
17use lenso_capability_secrets::{self as secrets, ResolveError, ResolveRequest, ResolveResponse};
18use lenso_kernel::RuntimeFailure;
19use secrecy::{ExposeSecret, SecretString};
20use zeroize::{Zeroize, Zeroizing};
21
22const SOURCE_PLACEHOLDER: &str = "{source}";
23const MAX_REFERENCE_LENGTH: usize = 256;
24const MAX_SOURCE_LENGTH: usize = 4096;
25const MAX_ARGUMENTS: usize = 64;
26const MAX_ARGUMENT_BYTES: usize = 65_536;
27const MAX_ENVIRONMENT_NAMES: usize = 64;
28const MAX_TIMEOUT_MS: u64 = 300_000;
29const MAX_OUTPUT_BYTES: usize = 1_048_576;
30
31#[derive(Clone, Debug, serde::Deserialize)]
32#[serde(deny_unknown_fields)]
33struct CommandConfig {
34    program: PathBuf,
35    arguments: Vec<String>,
36    environment_allowlist: Vec<String>,
37    #[serde(deserialize_with = "deserialize_unique_references")]
38    references: BTreeMap<String, String>,
39    timeout_ms: u64,
40    max_output_bytes: usize,
41}
42
43#[derive(Clone, Debug)]
44struct PreparedProgram {
45    executable: PathBuf,
46}
47
48fn validate_config(config: &CommandConfig) -> Result<(), RuntimeFailure> {
49    if !config.program.is_absolute() || config.program.as_os_str().is_empty() {
50        return Err(invalid_plan("Secrets command program must be absolute"));
51    }
52    if config.arguments.is_empty() || config.arguments.len() > MAX_ARGUMENTS {
53        return Err(invalid_plan(
54            "Secrets command arguments must contain between 1 and 64 entries",
55        ));
56    }
57    let argument_bytes = config
58        .arguments
59        .iter()
60        .try_fold(0usize, |total, argument| {
61            if argument.contains('\0') {
62                None
63            } else {
64                total.checked_add(argument.len())
65            }
66        })
67        .ok_or_else(|| invalid_plan("Secrets command arguments are invalid"))?;
68    if argument_bytes > MAX_ARGUMENT_BYTES
69        || config
70            .arguments
71            .iter()
72            .filter(|argument| argument.as_str() == SOURCE_PLACEHOLDER)
73            .count()
74            != 1
75        || config
76            .arguments
77            .iter()
78            .any(|argument| argument.contains(SOURCE_PLACEHOLDER) && argument != SOURCE_PLACEHOLDER)
79    {
80        return Err(invalid_plan(
81            "Secrets command arguments require exactly one standalone `{source}` placeholder",
82        ));
83    }
84    if config.environment_allowlist.len() > MAX_ENVIRONMENT_NAMES {
85        return Err(invalid_plan(
86            "Secrets command environment allowlist cannot exceed 64 names",
87        ));
88    }
89    let mut environment = BTreeSet::new();
90    if config
91        .environment_allowlist
92        .iter()
93        .any(|name| !valid_environment_variable(name) || !environment.insert(name))
94    {
95        return Err(invalid_plan(
96            "Secrets command environment allowlist contains invalid or duplicate names",
97        ));
98    }
99    if config.references.is_empty()
100        || config
101            .references
102            .iter()
103            .any(|(reference, source)| !valid_reference(reference) || !valid_source_name(source))
104    {
105        return Err(invalid_plan("Secrets command references are invalid"));
106    }
107    if !(1..=MAX_TIMEOUT_MS).contains(&config.timeout_ms) {
108        return Err(invalid_plan("timeout_ms must be between 1 and 300000"));
109    }
110    if !(1..=MAX_OUTPUT_BYTES).contains(&config.max_output_bytes) {
111        return Err(invalid_plan(
112            "max_output_bytes must be between 1 and 1048576",
113        ));
114    }
115    Ok(())
116}
117
118#[lenso::plugin(
119    lifecycle,
120    configuration_schema = "config.schema.json",
121    validate = validate_config
122)]
123#[derive(Clone, Debug)]
124struct CommandSecretsPlugin {
125    #[config]
126    config: CommandConfig,
127    prepared: Rc<RefCell<Option<PreparedProgram>>>,
128}
129
130impl Lifecycle for CommandSecretsPlugin {
131    async fn prepare(&self, _context: PrepareContext) -> Result<(), RuntimeFailure> {
132        let executable = fs::canonicalize(&self.config.program)
133            .map_err(|_| invalid_plan("Secrets command program is unavailable"))?;
134        if !executable.is_file() {
135            return Err(invalid_plan(
136                "Secrets command program is not a regular file",
137            ));
138        }
139        let prepared = PreparedProgram { executable };
140        self.prepared.replace(Some(prepared.clone()));
141        for (reference, source) in &self.config.references {
142            run_resolver(&self.config, &prepared, source)
143                .await
144                .map_err(|()| RuntimeFailure::PluginFailure {
145                    detail: format!(
146                        "configured command secret reference `{reference}` is unavailable"
147                    ),
148                })?;
149        }
150        Ok(())
151    }
152
153    fn deactivate(
154        &self,
155        _context: DeactivateContext,
156    ) -> impl std::future::Future<Output = Result<(), RuntimeFailure>> {
157        self.prepared.replace(None);
158        std::future::ready(Ok(()))
159    }
160}
161
162#[lenso::provides(secrets::Secrets)]
163impl CommandSecretsPlugin {
164    async fn resolve(
165        &self,
166        _context: Ctx,
167        request: ResolveRequest,
168    ) -> PluginResult<ResolveResponse, ResolveError> {
169        let ResolveRequest { reference } = request;
170        if !valid_reference(&reference) {
171            return Err(PluginError::domain(ResolveError::InvalidReference));
172        }
173        let source = self
174            .config
175            .references
176            .get(&reference)
177            .ok_or_else(|| PluginError::domain(ResolveError::UnknownReference))?;
178        let prepared = self.prepared.borrow().clone().ok_or_else(|| {
179            PluginError::runtime(RuntimeFailure::Unavailable {
180                capability: secrets::CAPABILITY_ID,
181            })
182        })?;
183        let value = run_resolver(&self.config, &prepared, source)
184            .await
185            .map_err(|()| {
186                PluginError::runtime(RuntimeFailure::PluginFailure {
187                    detail: format!(
188                        "configured command secret reference `{reference}` is unavailable"
189                    ),
190                })
191            })?;
192        Ok(ResolveResponse {
193            value: value.expose_secret().to_owned(),
194        })
195    }
196}
197
198async fn run_resolver(
199    config: &CommandConfig,
200    prepared: &PreparedProgram,
201    source: &str,
202) -> Result<SecretString, ()> {
203    let request = CommandRequest {
204        configured_program: config.program.clone(),
205        executable: prepared.executable.clone(),
206        arguments: config
207            .arguments
208            .iter()
209            .map(|argument| {
210                if argument == SOURCE_PLACEHOLDER {
211                    source.to_owned()
212                } else {
213                    argument.clone()
214                }
215            })
216            .collect(),
217        environment_allowlist: config.environment_allowlist.clone(),
218        timeout: Duration::from_millis(config.timeout_ms),
219        max_output_bytes: config.max_output_bytes,
220    };
221    let (sender, receiver) = oneshot::channel();
222    thread::Builder::new()
223        .name("lenso-secret-resolver".to_owned())
224        .spawn(move || {
225            let result = run_command(request, &sender);
226            let _ = sender.send(result);
227        })
228        .map_err(|_| ())?;
229    receiver.await.map_err(|_| ())?
230}
231
232#[derive(Debug)]
233struct CommandRequest {
234    configured_program: PathBuf,
235    executable: PathBuf,
236    arguments: Vec<String>,
237    environment_allowlist: Vec<String>,
238    timeout: Duration,
239    max_output_bytes: usize,
240}
241
242fn run_command(
243    request: CommandRequest,
244    cancellation: &oneshot::Sender<Result<SecretString, ()>>,
245) -> Result<SecretString, ()> {
246    let CommandRequest {
247        configured_program,
248        executable,
249        arguments,
250        environment_allowlist,
251        timeout,
252        max_output_bytes,
253    } = request;
254    if fs::canonicalize(&configured_program).map_err(|_| ())? != executable {
255        return Err(());
256    }
257    let mut command = Command::new(&executable);
258    command
259        .args(&arguments)
260        .env_clear()
261        .stdin(Stdio::null())
262        .stdout(Stdio::piped())
263        .stderr(Stdio::piped());
264    for name in &environment_allowlist {
265        if let Some(value) = env::var_os(name) {
266            command.env(name, value);
267        }
268    }
269    let mut child = command.spawn().map_err(|_| ())?;
270    let stdout = child.stdout.take().ok_or(())?;
271    let stderr = child.stderr.take().ok_or(())?;
272    let stdout_reader = read_bounded(stdout, max_output_bytes);
273    let stderr_reader = read_bounded(stderr, max_output_bytes);
274    let status = wait_for_child(&mut child, timeout, cancellation)?;
275    let stdout = stdout_reader.join().map_err(|_| ())??;
276    let _stderr = stderr_reader.join().map_err(|_| ())??;
277    if !status.success() || stdout.overflow {
278        return Err(());
279    }
280    let mut value = String::from_utf8(stdout.bytes.to_vec()).map_err(|_| ())?;
281    if value.ends_with('\n') {
282        value.pop();
283        if value.ends_with('\r') {
284            value.pop();
285        }
286    }
287    if value.is_empty() {
288        value.zeroize();
289        return Err(());
290    }
291    Ok(SecretString::from(value))
292}
293
294fn wait_for_child(
295    child: &mut Child,
296    timeout: Duration,
297    cancellation: &oneshot::Sender<Result<SecretString, ()>>,
298) -> Result<std::process::ExitStatus, ()> {
299    let started = Instant::now();
300    loop {
301        if let Some(status) = child.try_wait().map_err(|_| ())? {
302            return Ok(status);
303        }
304        if started.elapsed() >= timeout || cancellation.is_canceled() {
305            let _ = child.kill();
306            let _ = child.wait();
307            return Err(());
308        }
309        thread::sleep(Duration::from_millis(10));
310    }
311}
312
313#[derive(Debug)]
314struct BoundedOutput {
315    bytes: Zeroizing<Vec<u8>>,
316    overflow: bool,
317}
318
319fn read_bounded(
320    mut reader: impl Read + Send + 'static,
321    limit: usize,
322) -> thread::JoinHandle<Result<BoundedOutput, ()>> {
323    thread::spawn(move || {
324        let mut bytes = Zeroizing::new(Vec::new());
325        let mut buffer = [0u8; 8192];
326        let mut overflow = false;
327        loop {
328            let read = reader.read(&mut buffer).map_err(|_| ())?;
329            if read == 0 {
330                break;
331            }
332            let remaining = limit.saturating_sub(bytes.len());
333            let retained = remaining.min(read);
334            bytes.extend_from_slice(&buffer[..retained]);
335            overflow |= retained < read;
336            buffer[..read].zeroize();
337        }
338        buffer.zeroize();
339        Ok(BoundedOutput { bytes, overflow })
340    })
341}
342
343fn valid_reference(reference: &str) -> bool {
344    !reference.is_empty()
345        && reference.len() <= MAX_REFERENCE_LENGTH
346        && !reference.starts_with('/')
347        && !reference.ends_with('/')
348        && !reference.contains("//")
349        && !reference.contains('\0')
350        && reference
351            .split('/')
352            .all(|segment| segment != "." && segment != "..")
353}
354
355fn valid_source_name(value: &str) -> bool {
356    !value.trim().is_empty() && value.len() <= MAX_SOURCE_LENGTH && !value.contains('\0')
357}
358
359fn valid_environment_variable(value: &str) -> bool {
360    let mut bytes = value.bytes();
361    matches!(bytes.next(), Some(b'A'..=b'Z' | b'_'))
362        && bytes.all(|byte| byte.is_ascii_uppercase() || byte.is_ascii_digit() || byte == b'_')
363}
364
365fn deserialize_unique_references<'de, D>(
366    deserializer: D,
367) -> Result<BTreeMap<String, String>, D::Error>
368where
369    D: serde::Deserializer<'de>,
370{
371    struct UniqueReferences;
372
373    impl<'de> serde::de::Visitor<'de> for UniqueReferences {
374        type Value = BTreeMap<String, String>;
375
376        fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
377            formatter.write_str("a logical-reference to command-source map")
378        }
379
380        fn visit_map<A>(self, mut access: A) -> Result<Self::Value, A::Error>
381        where
382            A: serde::de::MapAccess<'de>,
383        {
384            let mut references = BTreeMap::new();
385            while let Some((reference, source)) = access.next_entry::<String, String>()? {
386                if references.insert(reference.clone(), source).is_some() {
387                    return Err(serde::de::Error::custom(format!(
388                        "duplicate logical secret reference `{reference}`"
389                    )));
390                }
391            }
392            Ok(references)
393        }
394    }
395
396    deserializer.deserialize_map(UniqueReferences)
397}
398
399fn invalid_plan(detail: impl Into<String>) -> RuntimeFailure {
400    RuntimeFailure::InvalidResolvedPlan {
401        detail: detail.into(),
402    }
403}
404
405#[cfg(all(test, unix))]
406mod tests {
407    use std::{io::Write, os::unix::fs::PermissionsExt, path::Path};
408
409    use super::*;
410
411    fn script(path: &Path, body: &str) {
412        let mut file = fs::File::create(path).unwrap();
413        writeln!(file, "#!/bin/sh").unwrap();
414        writeln!(file, "{body}").unwrap();
415        fs::set_permissions(path, fs::Permissions::from_mode(0o755)).unwrap();
416    }
417
418    fn config(program: PathBuf) -> CommandConfig {
419        CommandConfig {
420            program,
421            arguments: vec![SOURCE_PLACEHOLDER.to_owned()],
422            environment_allowlist: Vec::new(),
423            references: BTreeMap::from([(
424                "model/openai-api-key".to_owned(),
425                "op://vault/item/password".to_owned(),
426            )]),
427            timeout_ms: 2_000,
428            max_output_bytes: 4096,
429        }
430    }
431
432    #[test]
433    fn descriptor_exposes_one_secrets_provider() {
434        let descriptor: serde_json::Value = serde_json::from_str(PLUGIN_DESCRIPTOR_JSON).unwrap();
435        assert_eq!(descriptor["plugin_id"], "lenso.secrets.command");
436        assert_eq!(
437            descriptor["provided_capabilities"][0]["capability_id"],
438            "lenso.secrets@1"
439        );
440    }
441
442    #[test]
443    fn command_observes_rotation_and_strips_only_one_line_ending() {
444        let directory = tempfile::tempdir().unwrap();
445        let program = directory.path().join("resolver");
446        script(&program, r"printf 'first-secret\n'");
447        let config = config(program.clone());
448        let prepared = PreparedProgram {
449            executable: fs::canonicalize(&program).unwrap(),
450        };
451        let first = futures::executor::block_on(run_resolver(
452            &config,
453            &prepared,
454            "op://vault/item/password",
455        ))
456        .unwrap();
457        assert_eq!(first.expose_secret(), "first-secret");
458        assert!(!format!("{first:?}").contains("first-secret"));
459
460        script(&program, r"printf 'rotated-secret'");
461        let rotated = futures::executor::block_on(run_resolver(
462            &config,
463            &prepared,
464            "op://vault/item/password",
465        ))
466        .unwrap();
467        assert_eq!(rotated.expose_secret(), "rotated-secret");
468    }
469
470    #[test]
471    fn stderr_nonzero_timeout_and_output_overflow_are_redacted() {
472        let directory = tempfile::tempdir().unwrap();
473        let program = directory.path().join("resolver");
474        let mut config = config(program.clone());
475        script(&program, "echo never-log-this >&2; exit 7");
476        let prepared = PreparedProgram {
477            executable: fs::canonicalize(&program).unwrap(),
478        };
479        let failure = futures::executor::block_on(run_resolver(&config, &prepared, "source"));
480        assert!(failure.is_err());
481        assert!(!format!("{failure:?}").contains("never-log-this"));
482
483        script(&program, "sleep 1; printf secret");
484        config.timeout_ms = 10;
485        assert!(futures::executor::block_on(run_resolver(&config, &prepared, "source")).is_err());
486
487        script(&program, "printf 123456789");
488        config.timeout_ms = 2_000;
489        config.max_output_bytes = 4;
490        assert!(futures::executor::block_on(run_resolver(&config, &prepared, "source")).is_err());
491    }
492
493    #[test]
494    fn configuration_requires_absolute_program_and_one_standalone_placeholder() {
495        let mut invalid = config(PathBuf::from("relative"));
496        assert!(validate_config(&invalid).is_err());
497        invalid.program = PathBuf::from("/absolute/resolver");
498        invalid.arguments = vec!["prefix-{source}".to_owned()];
499        assert!(validate_config(&invalid).is_err());
500        invalid.arguments = vec![SOURCE_PLACEHOLDER.to_owned(), SOURCE_PLACEHOLDER.to_owned()];
501        assert!(validate_config(&invalid).is_err());
502    }
503}