chuggin 0.12.5

A terminal harness for persistent projects and continuous local model refinement
//! Owned command sessions. Yielding never kills a process; dropping its owner does.
use crate::{
    events::{self, Event},
    project::{Check, CheckResult},
};
use anyhow::{Context, Result};
use serde_json::{Value, json};
use std::{
    fs,
    io::{Read, Seek, SeekFrom, Write},
    path::{Path, PathBuf},
    process::{Child, Command, Stdio},
    sync::atomic::{AtomicBool, Ordering},
    thread,
    time::{Duration, Instant},
};

pub struct Session {
    pub id: String,
    pub log: PathBuf,
    pub argv: Vec<String>,
    child: Child,
    started: Instant,
    next_review: Instant,
    pub reviews: u64,
    pub result: Option<CheckResult>,
    pub stop_reason: Option<String>,
}
impl Session {
    pub fn start(root: &Path, c: &Check, log: &Path, id: &str) -> Result<Self> {
        anyhow::ensure!(!c.argv.is_empty(), "Command argv must not be empty");
        let file = fs::File::create(log)?;
        let mut cmd = Command::new(&c.argv[0]);
        #[cfg(unix)]
        {
            use std::os::unix::process::CommandExt;
            cmd.process_group(0);
        }
        let mut child = cmd
            .args(&c.argv[1..])
            .current_dir(root)
            .stdin(Stdio::piped())
            .stdout(file.try_clone()?)
            .stderr(file)
            .spawn()
            .context("Cannot start command")?;
        // Avoid blocking the agent if a process stops consuming stdin.
        #[cfg(unix)]
        {
            use std::os::fd::AsRawFd;
            if let Some(stdin) = &mut child.stdin {
                let fd = stdin.as_raw_fd();
                unsafe {
                    let flags = libc::fcntl(fd, libc::F_GETFL);
                    if flags >= 0 {
                        libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK);
                    }
                }
            }
        }
        crate::project::register_check(child.id() as i32);
        events::send(Event::Check {
            command: c.argv.join(" "),
            path: log.into(),
        });
        Ok(Self {
            id: id.into(),
            log: log.into(),
            argv: c.argv.clone(),
            child,
            started: Instant::now(),
            next_review: Instant::now() + Duration::from_secs(c.timeout_seconds.clamp(1, 86400)),
            reviews: 0,
            result: None,
            stop_reason: None,
        })
    }
    pub fn running(&self) -> bool {
        self.result.is_none()
    }
    pub fn due(&self) -> bool {
        self.running() && Instant::now() >= self.next_review
    }
    pub fn extend(&mut self, seconds: u64) {
        self.next_review = Instant::now() + Duration::from_secs(seconds.clamp(1, 900));
    }
    fn cleanup(&mut self) {
        #[cfg(unix)]
        unsafe {
            libc::kill(-(self.child.id() as i32), libc::SIGKILL);
        }
        let _ = self.child.kill();
        let _ = self.child.wait();
        crate::project::unregister_check(self.child.id() as i32);
    }
    pub fn tail(&self) -> Result<String> {
        let mut f = fs::File::open(&self.log)?;
        let len = f.metadata()?.len();
        f.seek(SeekFrom::Start(len.saturating_sub(7000)))?;
        let mut b = Vec::new();
        f.read_to_end(&mut b)?;
        Ok(String::from_utf8_lossy(&b).into())
    }
    fn finish(&mut self, exit_code: Option<i32>, passed: bool) -> Result<()> {
        self.cleanup();
        self.result = Some(CheckResult {
            exit_code,
            argv: self.argv.clone(),
            passed,
            timed_out: false,
            output: self.tail()?,
        });
        events::send(Event::CheckDone(passed));
        self.persist()?;
        Ok(())
    }
    pub fn poll(&mut self, milliseconds: u64, stop: &AtomicBool) -> Result<Value> {
        let deadline = Instant::now() + Duration::from_millis(milliseconds.min(1000));
        while self.running() {
            if let Some(status) = self.child.try_wait()? {
                self.finish(status.code(), status.success())?;
                break;
            }
            if stop.load(Ordering::SeqCst) {
                self.terminate("Stopped by operator")?;
                break;
            }
            if Instant::now() >= deadline {
                break;
            }
            thread::sleep(Duration::from_millis(25));
        }
        let v = self.snapshot()?;
        events::send(Event::CommandStatus {
            id: self.id.clone(),
            elapsed: self.started.elapsed().as_secs(),
            next_review: self
                .next_review
                .saturating_duration_since(Instant::now())
                .as_secs(),
            running: self.running(),
        });
        self.persist()?;
        Ok(v)
    }
    pub fn terminate(&mut self, reason: &str) -> Result<()> {
        // A process may have completed while its watchdog was considering the evidence.
        if !self.running() {
            return Ok(());
        }
        if let Some(status) = self.child.try_wait()? {
            return self.finish(status.code(), status.success());
        }
        self.stop_reason = Some(reason.into());
        self.finish(None, false)
    }
    pub fn input(&mut self, text: &str, close: bool) -> Result<usize> {
        anyhow::ensure!(self.running(), "Command has finished");
        anyhow::ensure!(text.len() <= 16000, "Input exceeds 16000 bytes");
        let stdin = self.child.stdin.as_mut().context("stdin is closed")?;
        let written = match stdin.write(text.as_bytes()) {
            Ok(n) => n,
            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => 0,
            Err(e) => return Err(e.into()),
        };
        if close && written == text.len() {
            self.child.stdin.take();
        }
        Ok(written)
    }
    pub fn snapshot(&self) -> Result<Value> {
        let log_id = format!(
            "{}/{}",
            self.log
                .parent()
                .and_then(Path::file_name)
                .unwrap_or_default()
                .to_string_lossy(),
            self.log.file_name().unwrap_or_default().to_string_lossy()
        );
        Ok(
            json!({"command_id":self.id,"argv":self.argv,"running":self.running(),"elapsed_seconds":self.started.elapsed().as_secs(),"next_review_seconds":self.next_review.saturating_duration_since(Instant::now()).as_secs(),"watchdog_reviews":self.reviews,"log_bytes":fs::metadata(&self.log)?.len(),"seconds_since_output":fs::metadata(&self.log)?.modified().ok().and_then(|t|t.elapsed().ok()).map(|d|d.as_secs()),"processes":process_observations(self.child.id()),"exit_code":self.result.as_ref().and_then(|r|r.exit_code),"passed":self.result.as_ref().map(|r|r.passed),"timed_out":false,"stop_reason":self.stop_reason,"output_tail":self.tail()?,"log_id":log_id,"instruction":"If running, inspect relevant evidence or use command_status to wait. Do not launch duplicates or edit files being checked. A running command is not a passing check."}),
        )
    }
    fn persist(&self) -> Result<()> {
        fs::write(
            self.log.with_extension("session.json"),
            serde_json::to_vec_pretty(&self.snapshot()?)?,
        )?;
        Ok(())
    }
}
impl Drop for Session {
    fn drop(&mut self) {
        if self.running() {
            let _ = self.terminate("Session owner ended; command interrupted, not verified");
        }
    }
}

/// Linux process-tree counters are evidence, not a hang heuristic. Other platforms
/// still supply elapsed time and log activity. Bounded traversal avoids huge output.
fn process_observations(pid: u32) -> Vec<Value> {
    let mut result = Vec::new();
    #[cfg(target_os = "linux")]
    {
        let mut pending = vec![pid];
        while let Some(id) = pending.pop() {
            if result.len() >= 24 {
                break;
            }
            if let Ok(stat) = fs::read_to_string(format!("/proc/{id}/stat")) {
                if let Some((_, fields)) = stat.rsplit_once(") ") {
                    let f: Vec<_> = fields.split_whitespace().collect();
                    result.push(json!({"pid":id,"state":f.first(),"user_cpu_ticks":f.get(11),"system_cpu_ticks":f.get(12),"resident_pages":f.get(21)}));
                }
                if let Ok(children) = fs::read_to_string(format!("/proc/{id}/task/{id}/children")) {
                    pending.extend(
                        children
                            .split_whitespace()
                            .take(24)
                            .filter_map(|p| p.parse::<u32>().ok()),
                    );
                }
            }
        }
    }
    #[cfg(not(target_os = "linux"))]
    {
        let _ = pid;
    }
    result
}

#[cfg(test)]
mod tests {
    use super::*;
    fn command(script: &str) -> Check {
        Check {
            argv: vec!["sh".into(), "-c".into(), script.into()],
            timeout_seconds: 1,
        }
    }
    #[test]
    fn yielding_preserves_process_and_stdin_and_completion() {
        let d = tempfile::tempdir().unwrap();
        let stop = AtomicBool::new(false);
        let mut s = Session::start(
            d.path(),
            &command("printf ready; read value; printf '\nreceived:%s' \"$value\""),
            &d.path().join("command-test.log"),
            "test",
        )
        .unwrap();
        assert_eq!(s.poll(100, &stop).unwrap()["running"], true);
        assert!(s.tail().unwrap().contains("ready"));
        s.input("hello\n", true).unwrap();
        assert_eq!(s.poll(1000, &stop).unwrap()["passed"], true);
        assert!(s.tail().unwrap().contains("received:hello"));
        assert_eq!(s.poll(0, &stop).unwrap()["passed"], true);
    }
    #[test]
    fn review_deadline_does_not_kill_and_stop_removes_descendants() {
        let d = tempfile::tempdir().unwrap();
        let stop = AtomicBool::new(false);
        let mut s = Session::start(
            d.path(),
            &command("(sleep 2; touch escaped) & wait"),
            &d.path().join("command-test.log"),
            "test",
        )
        .unwrap();
        s.next_review = Instant::now();
        assert!(s.due());
        assert_eq!(s.poll(0, &stop).unwrap()["running"], true);
        s.terminate("Test confirmed runaway").unwrap();
        assert_eq!(s.snapshot().unwrap()["passed"], false);
        thread::sleep(Duration::from_millis(2100));
        assert!(!d.path().join("escaped").exists());
    }
    #[test]
    fn completed_process_wins_over_stale_termination_and_drop_cleans_up() {
        let d = tempfile::tempdir().unwrap();
        let mut s = Session::start(
            d.path(),
            &command("exit 0"),
            &d.path().join("command-test.log"),
            "test",
        )
        .unwrap();
        thread::sleep(Duration::from_millis(100));
        s.terminate("stale diagnosis").unwrap();
        assert!(s.result.as_ref().unwrap().passed);
        assert!(s.stop_reason.is_none());
        {
            let _child = Session::start(
                d.path(),
                &command("sleep 1; touch escaped"),
                &d.path().join("command-drop.log"),
                "drop",
            )
            .unwrap();
        }
        thread::sleep(Duration::from_millis(1100));
        assert!(!d.path().join("escaped").exists());
    }
}