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);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Wakeup {
Mailbox {
id: String,
from: String,
to: String,
},
Event {
id: String,
name: Option<String>,
},
}
impl Wakeup {
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 },
}
}
}
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")
}
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);
}
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)
}
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()
);
}
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);
}
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, 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(),
}
);
}
}