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