1use std::path::{Path, PathBuf};
4#[cfg(unix)]
5use std::process::Stdio;
6use std::time::Duration;
7
8use mobius::backend::sandbox::ProcessGroupGuard;
9#[cfg(unix)]
10use nix::sys::signal::{Signal, kill};
11#[cfg(unix)]
12use nix::unistd::Pid;
13use serde::{Deserialize, Serialize};
14#[cfg(unix)]
15use tokio::process::Command;
16use tokio::process::{Child, ChildStdin};
17use tokio::task::JoinSet;
18use tokio::time::Instant;
19
20use crate::{Error, Result, host::GatewayHost};
21
22#[derive(Clone, Copy, Deserialize)]
23#[serde(deny_unknown_fields)]
24struct Defaults {
25 idle_grace_seconds: u64,
26 timeout_seconds: u64,
27 retry_seconds: u64,
28}
29mobius::embedded_config! {
30 static DEFAULTS: Defaults = include_str!("activity.toml");
31}
32
33#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
35#[serde(default, deny_unknown_fields)]
36pub struct ActivityHookConfig {
37 pub command: Vec<String>,
39 pub idle_grace_seconds: u64,
41 pub timeout_seconds: u64,
43 pub retry_seconds: u64,
45}
46
47impl Default for ActivityHookConfig {
48 fn default() -> Self {
49 Self {
50 command: Vec::new(),
51 idle_grace_seconds: DEFAULTS.idle_grace_seconds,
52 timeout_seconds: DEFAULTS.timeout_seconds,
53 retry_seconds: DEFAULTS.retry_seconds,
54 }
55 }
56}
57
58impl ActivityHookConfig {
59 pub fn validate(&self) -> Result<()> {
63 self.executable().map(drop)
64 }
65
66 pub fn validate_roots<'a>(&self, roots: impl IntoIterator<Item = &'a Path>) -> Result<()> {
70 let executable = self.executable()?;
71 for root in roots {
72 if executable.starts_with(std::fs::canonicalize(root)?) {
73 return Err(Error::Config(
74 "activity hook executable must be outside state and workspace roots".into(),
75 ));
76 }
77 }
78 Ok(())
79 }
80
81 fn executable(&self) -> Result<PathBuf> {
82 if !cfg!(unix) {
83 return Err(Error::Config("activity hooks require Unix".into()));
84 }
85 crate::config::bounded(
86 "telemetry.activity_hook.command",
87 self.command.len(),
88 1..=64,
89 )?;
90 for argument in &self.command {
91 if argument.len() > 4096 || argument.chars().any(char::is_control) {
92 return Err(Error::Config(
93 "activity hook arguments must be at most 4096 bytes without control characters"
94 .into(),
95 ));
96 }
97 }
98 crate::config::bounded(
99 "telemetry.activity_hook.idle_grace_seconds",
100 self.idle_grace_seconds,
101 0..=3600,
102 )?;
103 crate::config::bounded(
104 "telemetry.activity_hook.timeout_seconds",
105 self.timeout_seconds,
106 1..=60,
107 )?;
108 crate::config::bounded(
109 "telemetry.activity_hook.retry_seconds",
110 self.retry_seconds,
111 1..=3600,
112 )?;
113 let path = Path::new(&self.command[0]);
114 if !path.is_absolute() {
115 return Err(Error::Config(
116 "activity hook executable must be absolute".into(),
117 ));
118 }
119 let path = std::fs::canonicalize(path)?;
120 let metadata = path.metadata()?;
121 if !metadata.is_file() {
122 return Err(Error::Config(
123 "activity hook executable must be a file".into(),
124 ));
125 }
126 #[cfg(unix)]
127 {
128 use std::os::unix::fs::PermissionsExt as _;
129 if metadata.permissions().mode() & 0o111 == 0 {
130 return Err(Error::Config("activity hook file is not executable".into()));
131 }
132 }
133 Ok(path)
134 }
135}
136
137struct HookProcess {
138 child: Child,
139 group: ProcessGroupGuard,
140 stdin: Option<ChildStdin>,
142}
143
144pub(crate) struct ActivityHook<'a> {
145 config: Option<&'a ActivityHookConfig>,
146 executable: Option<PathBuf>,
147 process: Option<HookProcess>,
148 retiring: JoinSet<()>,
149}
150
151impl<'a> ActivityHook<'a> {
152 pub(crate) fn new(config: Option<&'a ActivityHookConfig>) -> Result<Self> {
153 Ok(Self {
154 executable: config.map(ActivityHookConfig::executable).transpose()?,
155 config,
156 process: None,
157 retiring: JoinSet::new(),
158 })
159 }
160
161 pub(crate) async fn run(&mut self, host: &GatewayHost, clients: impl Fn() -> Result<usize>) {
163 let Some(config) = self.config else {
164 return std::future::pending().await;
165 };
166 let retry = Duration::from_secs(config.retry_seconds);
167 let grace = Duration::from_secs(config.idle_grace_seconds);
168 let mut changes = host.activity_changes();
169 let mut revision = *changes.borrow_and_update();
170 let mut idle_at = None;
171 let mut retry_at = None;
172 self.acquire_when_due(&mut retry_at);
174 loop {
175 let current = *changes.borrow_and_update();
176 if current != revision {
177 revision = current;
178 idle_at = Some(Instant::now() + grace);
179 self.acquire_when_due(&mut retry_at);
180 }
181 let mut dirty = false;
182 let measured = {
183 let measurement = tokio::time::timeout(self.timeout(), host.runtime_activity());
184 tokio::pin!(measurement);
185 loop {
186 tokio::select! {
187 result = &mut measurement => break result,
188 result = wait_process(&mut self.process) => {
189 self.on_exit(result, &mut retry_at, retry);
190 }
191 () = tokio::time::sleep_until(retry_at.unwrap_or_else(Instant::now)),
192 if self.process.is_none() && retry_at.is_some() => {
193 self.acquire_when_due(&mut retry_at);
194 }
195 Ok(()) = changes.changed() => {
196 dirty = true;
197 let current = *changes.borrow_and_update();
198 if current != revision {
199 revision = current;
200 idle_at = Some(Instant::now() + grace);
201 self.acquire_when_due(&mut retry_at);
202 }
203 }
204 }
205 }
206 };
207 let mut idle = measured_idle(measured, &clients);
208 dirty |= changes.has_changed().unwrap_or(true);
209 let current = *changes.borrow_and_update();
210 if current != revision {
211 revision = current;
212 idle = false;
213 } else if dirty {
214 self.acquire_when_due(&mut retry_at);
216 continue;
217 }
218 let now = Instant::now();
219 if !idle {
220 idle_at = None;
221 }
222 let held = !idle || now < *idle_at.get_or_insert(now + grace);
223 if held {
224 self.acquire_when_due(&mut retry_at);
225 } else {
226 retry_at = None;
227 self.retire();
228 }
229 let mut next = Instant::now() + retry;
230 if held
231 && self.process.is_none()
232 && let Some(deadline) = retry_at
233 {
234 next = next.min(deadline);
235 }
236 if idle
237 && held
238 && let Some(deadline) = idle_at
239 {
240 next = next.min(deadline);
241 }
242 tokio::select! {
243 Ok(()) = changes.changed() => {}
244 () = tokio::time::sleep_until(next) => {}
245 result = wait_process(&mut self.process) => {
246 self.on_exit(result, &mut retry_at, retry);
247 }
248 Some(result) = self.retiring.join_next(), if !self.retiring.is_empty() => {
249 if let Err(error) = result {
250 tracing::warn!(%error, "activity hook cleanup failed");
251 }
252 }
253 }
254 }
255 }
256
257 fn acquire_when_due(&mut self, retry_at: &mut Option<Instant>) {
258 if self.process.is_none() && retry_at.is_none_or(|deadline| Instant::now() >= deadline) {
259 match self.acquire() {
260 Ok(()) => *retry_at = None,
261 Err(error) => {
262 tracing::warn!(%error, "activity hook acquisition failed");
263 *retry_at = Some(
264 Instant::now()
265 + Duration::from_secs(
266 self.config
267 .map_or(DEFAULTS.retry_seconds, |config| config.retry_seconds),
268 ),
269 );
270 }
271 }
272 }
273 }
274
275 fn on_exit(
276 &mut self,
277 result: std::io::Result<std::process::ExitStatus>,
278 retry_at: &mut Option<Instant>,
279 retry: Duration,
280 ) {
281 match result {
282 Ok(status) => tracing::warn!(%status, "activity hook command exited"),
283 Err(error) => tracing::warn!(%error, "activity hook command status failed"),
284 }
285 self.process = None;
286 *retry_at = Some(Instant::now() + retry);
287 }
288
289 fn retire(&mut self) {
290 if self.retiring.is_empty()
291 && let Some(process) = self.process.take()
292 {
293 self.retiring.spawn(stop_process(process, self.timeout()));
294 }
295 }
296
297 fn timeout(&self) -> Duration {
298 Duration::from_secs(
299 self.config
300 .map_or(DEFAULTS.timeout_seconds, |config| config.timeout_seconds),
301 )
302 }
303
304 #[cfg(unix)]
305 fn acquire(&mut self) -> Result<()> {
306 let Some(config) = self.config else {
307 return Ok(());
308 };
309 let Some(executable) = self.executable.as_ref() else {
310 return Ok(());
311 };
312 if self.process.is_some() {
313 return Ok(());
314 }
315 let mut child = Command::new(executable)
316 .args(&config.command[1..])
317 .env_clear()
318 .env("PATH", "/usr/bin:/bin")
319 .current_dir("/")
320 .stdin(Stdio::piped())
321 .stdout(Stdio::null())
322 .stderr(Stdio::inherit())
323 .process_group(0)
324 .kill_on_drop(true)
325 .spawn()?;
326 let group = ProcessGroupGuard::new(&child)?;
327 let stdin = child
328 .stdin
329 .take()
330 .ok_or_else(|| Error::Config("activity hook stdin unavailable".into()))?;
331 self.process = Some(HookProcess {
332 child,
333 group,
334 stdin: Some(stdin),
335 });
336 Ok(())
337 }
338
339 #[cfg(not(unix))]
340 fn acquire(&mut self) -> Result<()> {
341 Err(Error::Config("activity hooks require Unix".into()))
342 }
343
344 pub(crate) async fn stop(&mut self) {
345 if let Some(process) = self.process.take() {
346 stop_process(process, self.timeout()).await;
347 }
348 while let Some(result) = self.retiring.join_next().await {
349 if let Err(error) = result {
350 tracing::warn!(%error, "activity hook cleanup failed");
351 }
352 }
353 }
354}
355
356fn measured_idle(
357 measured: std::result::Result<
358 std::result::Result<crate::host::RuntimeActivity, crate::host::Rejection>,
359 tokio::time::error::Elapsed,
360 >,
361 clients: impl Fn() -> Result<usize>,
362) -> bool {
363 match measured {
364 Ok(Ok(activity)) => match clients() {
365 Ok(clients) => activity.idle && clients == 0,
366 Err(error) => {
367 tracing::warn!(%error, "activity hook client count failed");
368 false
369 }
370 },
371 Ok(Err(error)) => {
372 tracing::warn!(error = %error.message, "activity hook measurement failed");
373 false
374 }
375 Err(_) => {
376 tracing::warn!("activity hook measurement timed out");
377 false
378 }
379 }
380}
381
382async fn wait_process(
383 process: &mut Option<HookProcess>,
384) -> std::io::Result<std::process::ExitStatus> {
385 match process {
386 Some(process) => process.child.wait().await,
387 None => std::future::pending().await,
388 }
389}
390
391async fn stop_process(mut process: HookProcess, timeout: Duration) {
392 drop(process.stdin.take());
393 if !matches!(
394 tokio::time::timeout(timeout, process.child.wait()).await,
395 Ok(Ok(_))
396 ) {
397 #[cfg(unix)]
398 if let Some(pid) = process.child.id().and_then(|pid| i32::try_from(pid).ok()) {
399 let _ = kill(Pid::from_raw(-pid), Signal::SIGTERM);
400 }
401 if !matches!(
402 tokio::time::timeout(timeout, process.child.wait()).await,
403 Ok(Ok(_))
404 ) {
405 process.group.kill();
406 if !matches!(
407 tokio::time::timeout(timeout, process.child.wait()).await,
408 Ok(Ok(_))
409 ) {
410 tracing::warn!("activity hook command did not reap after termination");
411 }
412 }
413 }
414}
415
416#[cfg(all(test, unix))]
417mod tests {
418 use super::*;
419
420 fn command(arguments: &[&str]) -> ActivityHookConfig {
421 ActivityHookConfig {
422 command: arguments
423 .iter()
424 .map(|argument| (*argument).into())
425 .collect(),
426 timeout_seconds: 1,
427 ..Default::default()
428 }
429 }
430
431 #[test]
432 fn defaults_and_command_boundary_are_strict() {
433 let config: ActivityHookConfig = toml::from_str("command=['/bin/sh']").unwrap();
434 assert_eq!(
435 (
436 config.idle_grace_seconds,
437 config.timeout_seconds,
438 config.retry_seconds
439 ),
440 (5, 5, 5)
441 );
442 config.validate().unwrap();
443 assert!(
444 toml::from_str::<ActivityHookConfig>("command=['/bin/sh']\nretr_seconds=5").is_err()
445 );
446 for arguments in [&[][..], &["sh"][..], &["/bin/sh", "bad\nargument"][..]] {
447 assert!(command(arguments).validate().is_err());
448 }
449 let mut config = command(&["/bin/sh"]);
450 config.command.extend((0..64).map(|_| String::new()));
451 assert!(config.validate().is_err());
452 }
453
454 #[test]
455 fn configured_deadlines_are_bounded() {
456 let mut config = command(&["/bin/sh"]);
457 config.idle_grace_seconds = 3601;
458 assert!(config.validate().is_err());
459 config.idle_grace_seconds = 0;
460 for timeout in [0, 61] {
461 config.timeout_seconds = timeout;
462 assert!(config.validate().is_err());
463 }
464 config.timeout_seconds = 1;
465 for retry in [0, 3601] {
466 config.retry_seconds = retry;
467 assert!(config.validate().is_err());
468 }
469 config.retry_seconds = 1;
470 config.validate().unwrap();
471 }
472
473 #[test]
474 fn writable_roots_and_symlinks_cannot_supply_the_hook() {
475 let directory = tempfile::tempdir().unwrap();
476 let executable = directory.path().join("hook");
477 std::fs::write(&executable, "#!/bin/sh\nexit 0\n").unwrap();
478 std::fs::set_permissions(&executable, mobius::owner_only::file()).unwrap();
479 let config = command(&[executable.to_str().unwrap()]);
480 assert!(config.validate().is_err());
481 use std::os::unix::fs::PermissionsExt as _;
482 std::fs::set_permissions(&executable, std::fs::Permissions::from_mode(0o700)).unwrap();
483 config.validate().unwrap();
484 assert!(config.validate_roots([directory.path()]).is_err());
485 let link = directory.path().join("alias");
486 std::os::unix::fs::symlink(&executable, &link).unwrap();
487 assert!(
488 command(&[link.to_str().unwrap()])
489 .validate_roots([directory.path()])
490 .is_err()
491 );
492 }
493
494 #[tokio::test]
495 async fn idle_release_closes_stdin_for_foreground_cleanup() {
496 let directory = tempfile::tempdir().unwrap();
497 let marker = directory.path().join("released");
498 let config = command(&[
499 "/bin/sh",
500 "-c",
501 "cat >/dev/null; printf eof > \"$1\"",
502 "hook",
503 marker.to_str().unwrap(),
504 ]);
505 let mut hook = ActivityHook::new(Some(&config)).unwrap();
506 hook.acquire().unwrap();
507 assert!(hook.process.is_some());
508 hook.stop().await;
509 assert_eq!(std::fs::read_to_string(marker).unwrap(), "eof");
510 assert!(hook.process.is_none());
511 }
512
513 #[tokio::test]
514 async fn uncooperative_hook_is_killed_and_reaped_within_the_deadline() {
515 let config = command(&["/bin/sh", "-c", "trap '' TERM; while :; do sleep 1; done"]);
516 let mut hook = ActivityHook::new(Some(&config)).unwrap();
517 hook.acquire().unwrap();
518 let pid = i32::try_from(hook.process.as_ref().unwrap().child.id().unwrap()).unwrap();
519 tokio::time::timeout(Duration::from_secs(4), hook.stop())
520 .await
521 .unwrap();
522 assert_eq!(
523 kill(Pid::from_raw(pid), None),
524 Err(nix::errno::Errno::ESRCH)
525 );
526 }
527
528 #[tokio::test]
529 async fn new_work_reacquires_while_the_previous_hook_retires() {
530 use std::cell::Cell;
531
532 let directory = tempfile::tempdir().unwrap();
533 let (host, _bots) = empty_gateway(directory.path()).await;
534 let marker = directory.path().join("leases");
535 let mut config = command(&[
536 "/bin/sh",
537 "-c",
538 "printf 'start\\n' >> \"$1\"; cat >/dev/null; printf 'retire\\n' >> \"$1\"; sleep 3",
539 "hook",
540 marker.to_str().unwrap(),
541 ]);
542 config.idle_grace_seconds = 0;
543 config.timeout_seconds = 5;
544 let clients = Cell::new(1);
545 let mut hook = ActivityHook::new(Some(&config)).unwrap();
546 {
547 let activity = hook.run(&host, || Ok(clients.get()));
548 tokio::pin!(activity);
549 tokio::select! {
550 () = &mut activity => panic!("activity controller stopped"),
551 () = async {
552 wait_for_marker(&marker, "start", 1, Duration::from_secs(1)).await;
553 clients.set(0);
554 host.mark_runtime_activity();
555 wait_for_marker(&marker, "retire", 1, Duration::from_secs(1)).await;
556 clients.set(1);
557 host.mark_runtime_activity();
558 wait_for_marker(&marker, "start", 2, Duration::from_secs(1)).await;
560 } => {}
561 }
562 }
563 hook.stop().await;
564 host.shutdown().await;
565 }
566
567 #[tokio::test]
568 async fn a_late_storage_wake_requires_a_fresh_measurement_before_release() {
569 use std::cell::Cell;
570
571 let directory = tempfile::tempdir().unwrap();
572 let (host, bots) = empty_gateway(directory.path()).await;
573 let mut config = command(&["/bin/cat"]);
574 config.idle_grace_seconds = 0;
575 let measurements = Cell::new(0);
576 let mut hook = ActivityHook::new(Some(&config)).unwrap();
577 {
578 let activity = hook.run(&host, || {
579 measurements.set(measurements.get() + 1);
580 if measurements.get() == 1 {
581 bots.create_bot("Changed Bot", "A committed change.", Default::default())?;
583 Ok(0)
584 } else {
585 Ok(1)
586 }
587 });
588 tokio::pin!(activity);
589 tokio::select! {
590 () = &mut activity => panic!("activity controller stopped"),
591 result = tokio::time::timeout(Duration::from_secs(1), async {
592 while measurements.get() < 2 {
593 tokio::time::sleep(Duration::from_millis(10)).await;
594 }
595 }) => result.expect("late storage wake was measured again promptly"),
596 }
597 }
598 assert!(
599 hook.process.is_some(),
600 "the hold survives the stale idle measurement"
601 );
602 hook.stop().await;
603 host.shutdown().await;
604 }
605
606 #[tokio::test]
607 async fn unexpected_command_exit_retries_without_an_extra_health_poll_delay() {
608 let directory = tempfile::tempdir().unwrap();
609 let (host, _bots) = empty_gateway(directory.path()).await;
610 let marker = directory.path().join("starts");
611 let mut config = command(&[
612 "/bin/sh",
613 "-c",
614 "printf 'start\\n' >> \"$1\"; sleep 0.1",
615 "hook",
616 marker.to_str().unwrap(),
617 ]);
618 config.retry_seconds = 1;
619 let mut hook = ActivityHook::new(Some(&config)).unwrap();
620 {
621 let activity = hook.run(&host, || Ok(1));
622 tokio::pin!(activity);
623 tokio::select! {
624 () = &mut activity => panic!("activity controller stopped"),
625 () = async {
626 wait_for_marker(&marker, "start", 1, Duration::from_secs(1)).await;
627 wait_for_marker(&marker, "start", 2, Duration::from_millis(1800)).await;
628 } => {}
629 }
630 }
631 hook.stop().await;
632 host.shutdown().await;
633 }
634
635 async fn empty_gateway(path: &Path) -> (GatewayHost, std::sync::Arc<crate::bots::BotStore>) {
636 use crate::bots::BotStore;
637 use crate::config::{ConfigStore, CredentialStore};
638 use std::sync::Arc;
639
640 let (store, config) =
641 ConfigStore::initialize(path.join("state"), "127.0.0.1:8741".parse().unwrap(), None)
642 .unwrap();
643 let credentials = Arc::new(CredentialStore::open(store.credentials_path()).unwrap());
644 let bots = Arc::new(BotStore::open(store.state_dir()).unwrap());
645 let host = GatewayHost::start(store, config, credentials, Arc::clone(&bots))
646 .await
647 .unwrap();
648 (host, bots)
649 }
650
651 async fn wait_for_marker(path: &Path, value: &str, count: usize, timeout: Duration) {
652 tokio::time::timeout(timeout, async {
653 loop {
654 let text = tokio::fs::read_to_string(path).await.unwrap_or_default();
655 if text.lines().filter(|line| *line == value).count() >= count {
656 break;
657 }
658 tokio::time::sleep(Duration::from_millis(10)).await;
659 }
660 })
661 .await
662 .expect("activity lease transition was prompt");
663 }
664}