jan-cli 0.24.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.
//!
//! Tick-phase notices (`kind: tick`) are counted but not turned into agent wakeups;
//! Jan's tick driver (separate) owns Unifier tick lifecycles via `tick.sock`.

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, AtomicU64, 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);

/// Counters for the events.sock listener (shared with cron status).
#[derive(Debug, Default)]
pub struct WakeupMetrics {
    /// Mailbox/event notices enqueued for spawn.
    pub enqueued: AtomicU64,
    /// Dropped because the wakeup queue was full (oldest discarded).
    pub dropped: AtomicU64,
    /// `kind: tick` notices observed (not enqueued as agent wakeups).
    pub tick_notices: AtomicU64,
    /// Unparseable / unsupported notice lines.
    pub ignored: AtomicU64,
    /// Successful (re)connections to events.sock.
    pub reconnects: AtomicU64,
}

impl WakeupMetrics {
    pub fn snapshot(&self) -> WakeupMetricsSnapshot {
        WakeupMetricsSnapshot {
            enqueued: self.enqueued.load(Ordering::Relaxed),
            dropped: self.dropped.load(Ordering::Relaxed),
            tick_notices: self.tick_notices.load(Ordering::Relaxed),
            ignored: self.ignored.load(Ordering::Relaxed),
            reconnects: self.reconnects.load(Ordering::Relaxed),
        }
    }
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WakeupMetricsSnapshot {
    pub enqueued: u64,
    pub dropped: u64,
    pub tick_notices: u64,
    pub ignored: u64,
    pub reconnects: u64,
}

/// 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>,
    },
    /// ACID tick phase broadcast — not an agent wakeup.
    Tick {
        #[allow(dead_code)]
        tick: u64,
        #[allow(dead_code)]
        phase: String,
        #[serde(default)]
        #[allow(dead_code)]
        label: Option<String>,
    },
}

impl Notice {
    fn into_wakeup(self) -> Option<Wakeup> {
        match self {
            Self::Mailbox { id, from, to } => Some(Wakeup::Mailbox { id, from, to }),
            Self::Event { id, name } => Some(Wakeup::Event { id, name }),
            Self::Tick { .. } => None,
        }
    }

    fn is_tick(&self) -> bool {
        matches!(self, Self::Tick { .. })
    }
}

/// 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")
}

fn classify_notice_line(line: &str) -> Result<Notice, ()> {
    let trimmed = line.trim();
    if trimmed.is_empty() {
        return Err(());
    }
    serde_json::from_str::<Notice>(trimmed).map_err(|_| ())
}

fn push_wakeup(
    queue: &Mutex<VecDeque<Wakeup>>,
    metrics: &WakeupMetrics,
    wakeup: Wakeup,
) {
    let Ok(mut q) = queue.lock() else {
        return;
    };
    if q.len() >= MAX_WAKEUPS {
        q.pop_front();
        metrics.dropped.fetch_add(1, Ordering::Relaxed);
    }
    q.push_back(wakeup);
    metrics.enqueued.fetch_add(1, Ordering::Relaxed);
}

/// 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>,
    metrics: Arc<WakeupMetrics>,
    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);
                metrics.reconnects.fetch_add(1, Ordering::Relaxed);
                if verbose {
                    eprintln!(
                        "jan cron daemon: connected to unifier events.sock at {}",
                        path.display()
                    );
                }
                if !read_loop(&stream, &queue, &metrics, 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>>,
    metrics: &WakeupMetrics,
    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(_) => match classify_notice_line(&line) {
                Ok(notice) if notice.is_tick() => {
                    metrics.tick_notices.fetch_add(1, Ordering::Relaxed);
                    if verbose {
                        eprintln!("jan cron daemon: tick notice (not an agent wakeup): {}", line.trim());
                    }
                }
                Ok(notice) => {
                    if let Some(wakeup) = notice.into_wakeup() {
                        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, metrics, wakeup);
                    }
                }
                Err(()) => {
                    let preview = line.trim();
                    if !preview.is_empty() {
                        metrics.ignored.fetch_add(1, Ordering::Relaxed);
                        if verbose {
                            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::*;

    fn parse_notice_line(line: &str) -> Option<Wakeup> {
        classify_notice_line(line).ok().and_then(Notice::into_wakeup)
    }

    #[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 tick_notice_is_not_an_agent_wakeup() {
        let line = r#"{"kind":"tick","tick":2,"phase":"sense","label":"smoke"}"#;
        assert!(parse_notice_line(line).is_none());
        let notice = classify_notice_line(line).unwrap();
        assert!(notice.is_tick());
    }

    #[test]
    fn queue_caps_at_max_and_counts_drops() {
        let q = Mutex::new(VecDeque::new());
        let metrics = WakeupMetrics::default();
        for i in 0..(MAX_WAKEUPS + 10) {
            push_wakeup(
                &q,
                &metrics,
                Wakeup::Mailbox {
                    id: format!("{i}"),
                    from: "a".into(),
                    to: "b".into(),
                },
            );
        }
        assert_eq!(wakeup_depth(&q), MAX_WAKEUPS);
        assert_eq!(metrics.dropped.load(Ordering::Relaxed), 10);
        assert_eq!(metrics.enqueued.load(Ordering::Relaxed), MAX_WAKEUPS as u64 + 10);
        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(),
            }
        );
    }
}