1use 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}