use std::fs::File;
use std::io::{Error, Result};
use std::os::unix::fs::{DirBuilderExt, MetadataExt, OpenOptionsExt, PermissionsExt};
use std::os::unix::io::AsRawFd;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Mutex, OnceLock};
use std::time::Duration;
use tokio::net::{TcpListener, UnixListener};
use crate::router::{self, Registry};
use crate::{door, pair};
pub fn runtime_dir() -> PathBuf {
match std::env::var_os("XDG_RUNTIME_DIR") {
Some(r) => PathBuf::from(r).join("zap"),
None => directories::BaseDirs::new()
.map(|d| d.home_dir().to_path_buf())
.unwrap_or_else(|| PathBuf::from("."))
.join(".zap")
.join("run"),
}
}
pub fn socket_path() -> PathBuf {
runtime_dir().join("zapd.sock")
}
pub(crate) fn private_runtime() -> Result<PathBuf> {
let dir = runtime_dir();
std::fs::DirBuilder::new()
.recursive(true)
.mode(0o700)
.create(&dir)?;
let m = std::fs::symlink_metadata(&dir)?;
let me = unsafe { libc::getuid() };
if !m.file_type().is_dir() || m.uid() != me {
return Err(Error::new(
std::io::ErrorKind::PermissionDenied,
format!(
"zapd: {} must be a directory owned by uid {me}",
dir.display()
),
));
}
if m.mode() & 0o077 != 0 {
tracing::warn!("zapd: {} was open to others; made it 0700", dir.display());
std::fs::set_permissions(&dir, std::fs::Permissions::from_mode(0o700))?;
}
Ok(dir)
}
pub fn embed() {
static STARTED: OnceLock<()> = OnceLock::new();
STARTED.get_or_init(|| {
std::thread::Builder::new()
.name("zapd".into())
.spawn(|| {
let rt = crate::runtime();
loop {
if let Err(e) = lock(libc::F_SETLKW, true) {
tracing::error!("zapd: cannot take the router lock: {e}");
std::thread::sleep(Duration::from_secs(1));
continue;
}
ELECTED.store(true, Ordering::SeqCst);
tracing::info!("zapd: elected router (pid {})", std::process::id());
if let Err(e) = rt.block_on(serve()) {
tracing::error!("zapd: router failed: {e}");
}
ELECTED.store(false, Ordering::SeqCst);
let _ = lock(libc::F_SETLK, false);
std::thread::sleep(Duration::from_secs(1));
}
})
.expect("spawn the zapd thread");
});
}
static ELECTED: AtomicBool = AtomicBool::new(false);
pub fn holder() -> Option<u32> {
if ELECTED.load(Ordering::SeqCst) {
return Some(std::process::id());
}
let mut l: libc::flock = unsafe { std::mem::zeroed() };
l.l_type = libc::F_WRLCK as _;
l.l_whence = libc::SEEK_SET as _;
let fd = lock_fd().ok()?;
let ok = unsafe { libc::fcntl(fd, libc::F_GETLK, &mut l) } == 0;
let unlocked: libc::c_short = libc::F_UNLCK as _;
(ok && l.l_type != unlocked).then_some(l.l_pid as u32)
}
fn lock_fd() -> Result<std::os::unix::io::RawFd> {
static FILES: Mutex<Vec<(PathBuf, File)>> = Mutex::new(Vec::new());
let path = runtime_dir().join("zapd.lock");
let mut files = FILES.lock().unwrap();
if let Some((_, f)) = files.iter().find(|(p, _)| *p == path) {
return Ok(f.as_raw_fd());
}
private_runtime()?;
let f = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.mode(0o600)
.custom_flags(libc::O_NOFOLLOW)
.open(&path)?;
let fd = f.as_raw_fd();
files.push((path, f));
Ok(fd)
}
fn lock(cmd: libc::c_int, take: bool) -> Result<()> {
let fd = lock_fd()?;
let mut l: libc::flock = unsafe { std::mem::zeroed() };
l.l_type = if take { libc::F_WRLCK } else { libc::F_UNLCK } as _;
l.l_whence = libc::SEEK_SET as _;
loop {
if unsafe { libc::fcntl(fd, cmd, &l) } == 0 {
return Ok(());
}
let e = Error::last_os_error();
if e.raw_os_error() != Some(libc::EINTR) {
return Err(e);
}
}
}
async fn serve() -> Result<()> {
let path = private_runtime()?.join("zapd.sock");
let _ = std::fs::remove_file(&path); let uds = UnixListener::bind(&path)?;
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600))?;
tracing::info!("zapd: listening on {}", path.display());
let registry = Registry::new(&crate::host());
tokio::spawn(browser_door(registry.clone()));
loop {
match uds.accept().await {
Ok((stream, _)) => {
let registry = registry.clone();
tokio::spawn(async move {
if let Err(e) = router::handle(stream, registry, true).await {
tracing::debug!("zapd: connection ended: {e}");
}
});
}
Err(e) => tracing::warn!("zapd: accept: {e}"),
}
}
}
async fn browser_door(registry: Registry) {
let mut wait = Duration::from_millis(50);
let mut warned = false;
loop {
let port = match pair::load() {
Ok(p) => p.port,
Err(e) => {
tracing::error!("zapd: no pairing at {}: {e}", pair::path().display());
tokio::time::sleep(Duration::from_secs(5)).await;
continue;
}
};
match TcpListener::bind(("127.0.0.1", port)).await {
Ok(listener) => {
tracing::info!("zapd: browser door on 127.0.0.1:{port}");
door::serve(listener, port, registry).await;
return;
}
Err(e) => {
if !warned && wait >= Duration::from_secs(1) {
tracing::error!(
"zapd: the browser door 127.0.0.1:{port} is held by another program ({e}); browsers cannot connect until it frees. `zapd pair --reset` moves the door to a free port, and every browser pairs again"
);
warned = true;
}
tokio::time::sleep(wait).await;
wait = (wait * 2).min(Duration::from_secs(2));
}
}
}
}