Skip to main content

mobius_gateway/telemetry/
activity.rs

1//! Operator-owned host activity lease, independent of outbound collectors.
2
3use 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/// Local foreground command holding a host inhibitor until stdin closes.
34#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
35#[serde(default, deny_unknown_fields)]
36pub struct ActivityHookConfig {
37    /// Absolute trusted executable followed by its arguments; no shell is added.
38    pub command: Vec<String>,
39    /// Hold after the latest work transition or first confirmed idle observation.
40    pub idle_grace_seconds: u64,
41    /// Deadline for activity measurements and each graceful child shutdown phase.
42    pub timeout_seconds: u64,
43    /// Retry and child-health interval after a failed measurement or command exit.
44    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    /// Validates local command text, deadlines, and executable identity.
60    /// # Errors
61    /// Returns an error for invalid settings or an unavailable executable.
62    pub fn validate(&self) -> Result<()> {
63        self.executable().map(drop)
64    }
65
66    /// Rejects hook executables located inside protected or agent-writable roots.
67    /// # Errors
68    /// Returns an error for inaccessible roots or an overlapping executable.
69    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    // Child::wait closes Child.stdin, so the hold owns its write end separately.
141    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    /// Runs as a pinned server future; cancellation leaves cleanup owned by this hook.
162    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        // Unknown initial activity must not release the host before measurement succeeds.
173        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                // A storage wake can commit pending work without changing the execution revision.
215                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                        eprintln!("activity hook cleanup failed: {error}");
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                    eprintln!("activity hook acquisition failed: {error}");
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) => eprintln!("activity hook command exited ({status})"),
283            Err(error) => eprintln!("activity hook command status failed: {error}"),
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                eprintln!("activity hook cleanup failed: {error}");
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                eprintln!("activity hook client count failed: {error}");
368                false
369            }
370        },
371        Ok(Err(error)) => {
372            eprintln!("activity hook measurement failed: {}", error.message);
373            false
374        }
375        Err(_) => {
376            eprintln!("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                eprintln!("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                    // Cleanup sleeps three seconds; a new hold must precede its completion.
559                    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                    // The idle query already completed; commit a same-revision storage wake.
582                    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}