jan-cli 0.21.0

YAML-defined CLI trees with progressive help, optional exec aliases, merged extra specs, and SQLite audit logging keyed by git branch
Documentation
//! Listen on Unifier's `.daemon/events.sock` for mailbox/event wakeups.
//!
//! Unifier accepts subscribers on that socket and broadcasts line-delimited JSON
//! [`Notice`] values when agents send mail or post named events. Jan's cron
//! daemon keeps a durable connection (with reconnect) and queues wakeups so the
//! next 100 ms tick can spawn the matching script leaf.

use std::collections::VecDeque;
use std::io::{BufRead, BufReader, ErrorKind};
use std::os::unix::net::UnixStream;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::Duration;

use serde::Deserialize;

const MAX_WAKEUPS: usize = 256;
const RECONNECT_WAIT: Duration = Duration::from_secs(1);

/// Queued wakeup derived from a Unifier notice.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Wakeup {
    Mailbox {
        id: String,
        from: String,
        to: String,
    },
    Event {
        id: String,
        name: Option<String>,
    },
}

impl Wakeup {
    /// Script leaf name Jan should invoke, if any.
    pub fn agent_name(&self) -> Option<&str> {
        match self {
            Self::Mailbox { to, .. } => Some(to.as_str()),
            Self::Event { name, .. } => name.as_deref(),
        }
    }
}

#[derive(Debug, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
enum Notice {
    Mailbox {
        id: String,
        from: String,
        to: String,
    },
    Event {
        id: String,
        #[serde(default)]
        name: Option<String>,
    },
}

impl Notice {
    fn into_wakeup(self) -> Wakeup {
        match self {
            Self::Mailbox { id, from, to } => Wakeup::Mailbox { id, from, to },
            Self::Event { id, name } => Wakeup::Event { id, name },
        }
    }
}

/// Resolve Unifier's events socket (`$UNIFIER_HOME` or `~/.local/unifier`).
pub fn events_socket_path() -> PathBuf {
    unifier_home().join(".daemon").join("events.sock")
}

pub fn unifier_home() -> PathBuf {
    if let Ok(home) = std::env::var("UNIFIER_HOME") {
        let home = home.trim();
        if !home.is_empty() {
            return PathBuf::from(home);
        }
    }
    dirs::home_dir()
        .unwrap_or_else(|| PathBuf::from("."))
        .join(".local")
        .join("unifier")
}

/// Parse one line of Unifier event-socket JSON.
pub fn parse_notice_line(line: &str) -> Option<Wakeup> {
    let trimmed = line.trim();
    if trimmed.is_empty() {
        return None;
    }
    serde_json::from_str::<Notice>(trimmed)
        .ok()
        .map(Notice::into_wakeup)
}

fn push_wakeup(queue: &Mutex<VecDeque<Wakeup>>, wakeup: Wakeup) {
    let Ok(mut q) = queue.lock() else {
        return;
    };
    if q.len() >= MAX_WAKEUPS {
        q.pop_front();
    }
    q.push_back(wakeup);
}

/// Drain all pending wakeups (tick path).
pub fn drain_wakeups(queue: &Mutex<VecDeque<Wakeup>>) -> Vec<Wakeup> {
    let Ok(mut q) = queue.lock() else {
        return Vec::new();
    };
    q.drain(..).collect()
}

pub fn wakeup_depth(queue: &Mutex<VecDeque<Wakeup>>) -> usize {
    queue.lock().map(|q| q.len()).unwrap_or(0)
}

/// Background loop: connect to Unifier events.sock, enqueue notices, reconnect.
pub fn run_listener(
    queue: Arc<Mutex<VecDeque<Wakeup>>>,
    connected: Arc<AtomicBool>,
    verbose: bool,
) {
    let path = events_socket_path();
    if verbose {
        eprintln!(
            "jan cron daemon: unifier events listener targeting {}",
            path.display()
        );
    }

    while !crate::cron_daemon::stop_requested() {
        connected.store(false, Ordering::SeqCst);
        match UnixStream::connect(&path) {
            Ok(stream) => {
                if let Err(e) = stream.set_read_timeout(Some(Duration::from_millis(500))) {
                    if verbose {
                        eprintln!("jan cron daemon: events.sock set_read_timeout: {e:#}");
                    }
                    thread::sleep(RECONNECT_WAIT);
                    continue;
                }
                connected.store(true, Ordering::SeqCst);
                if verbose {
                    eprintln!(
                        "jan cron daemon: connected to unifier events.sock at {}",
                        path.display()
                    );
                }
                if !read_loop(&stream, &queue, verbose) {
                    break;
                }
                connected.store(false, Ordering::SeqCst);
                if verbose {
                    eprintln!("jan cron daemon: unifier events.sock disconnected; reconnecting");
                }
            }
            Err(e) => {
                if verbose {
                    eprintln!(
                        "jan cron daemon: waiting for unifier events.sock at {}: {e}",
                        path.display()
                    );
                }
                // Sleep in short slices so stop is responsive.
                let mut waited = Duration::ZERO;
                while waited < RECONNECT_WAIT && !crate::cron_daemon::stop_requested() {
                    thread::sleep(Duration::from_millis(100));
                    waited += Duration::from_millis(100);
                }
            }
        }
    }
    connected.store(false, Ordering::SeqCst);
}

/// Returns false when the global stop flag is set.
fn read_loop(stream: &UnixStream, queue: &Mutex<VecDeque<Wakeup>>, verbose: bool) -> bool {
    let mut reader = BufReader::new(stream);
    while !crate::cron_daemon::stop_requested() {
        let mut line = String::new();
        match reader.read_line(&mut line) {
            Ok(0) => return true, // peer closed; reconnect
            Ok(_) => {
                if let Some(wakeup) = parse_notice_line(&line) {
                    if verbose {
                        match &wakeup {
                            Wakeup::Mailbox { id, from, to } => {
                                eprintln!(
                                    "jan cron daemon: mailbox wakeup to={to} from={from} id={id}"
                                );
                            }
                            Wakeup::Event { id, name } => {
                                eprintln!(
                                    "jan cron daemon: event wakeup name={} id={id}",
                                    name.as_deref().unwrap_or("(none)")
                                );
                            }
                        }
                    }
                    push_wakeup(queue, wakeup);
                } else if verbose {
                    let preview = line.trim();
                    if !preview.is_empty() {
                        eprintln!("jan cron daemon: ignore unifier notice: {preview}");
                    }
                }
            }
            Err(e)
                if matches!(
                    e.kind(),
                    ErrorKind::Interrupted | ErrorKind::WouldBlock | ErrorKind::TimedOut
                ) =>
            {
                continue;
            }
            Err(e) => {
                if verbose {
                    eprintln!("jan cron daemon: events.sock read error: {e:#}");
                }
                return true;
            }
        }
    }
    false
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn parse_mailbox_notice() {
        let line = r#"{"kind":"mailbox","id":"11111111-1111-1111-1111-111111111111","from":"ping-agent","to":"pong-agent"}"#;
        let w = parse_notice_line(line).unwrap();
        assert_eq!(
            w,
            Wakeup::Mailbox {
                id: "11111111-1111-1111-1111-111111111111".into(),
                from: "ping-agent".into(),
                to: "pong-agent".into(),
            }
        );
        assert_eq!(w.agent_name(), Some("pong-agent"));
    }

    #[test]
    fn parse_event_notice() {
        let line = r#"{"kind":"event","id":"22222222-2222-2222-2222-222222222222","name":"status-report"}"#;
        let w = parse_notice_line(line).unwrap();
        assert_eq!(
            w,
            Wakeup::Event {
                id: "22222222-2222-2222-2222-222222222222".into(),
                name: Some("status-report".into()),
            }
        );
        assert_eq!(w.agent_name(), Some("status-report"));
    }

    #[test]
    fn parse_event_without_name() {
        let line = r#"{"kind":"event","id":"22222222-2222-2222-2222-222222222222"}"#;
        let w = parse_notice_line(line).unwrap();
        assert_eq!(w.agent_name(), None);
    }

    #[test]
    fn queue_caps_at_max() {
        let q = Mutex::new(VecDeque::new());
        for i in 0..(MAX_WAKEUPS + 10) {
            push_wakeup(
                &q,
                Wakeup::Mailbox {
                    id: format!("{i}"),
                    from: "a".into(),
                    to: "b".into(),
                },
            );
        }
        assert_eq!(wakeup_depth(&q), MAX_WAKEUPS);
        let drained = drain_wakeups(&q);
        assert_eq!(drained.len(), MAX_WAKEUPS);
        assert_eq!(
            drained[0],
            Wakeup::Mailbox {
                id: "10".into(),
                from: "a".into(),
                to: "b".into(),
            }
        );
    }
}