use std::os::fd::{AsRawFd, FromRawFd, OwnedFd};
use std::path::Path;
#[derive(Debug, PartialEq)]
pub(super) enum Woke {
Changed,
Rung,
TimedOut,
}
pub(super) struct Bell {
fd: Option<OwnedFd>,
}
impl Default for Bell {
fn default() -> Bell {
let fd = unsafe { libc::eventfd(0, libc::EFD_NONBLOCK | libc::EFD_CLOEXEC) };
let fd = (fd >= 0).then(|| unsafe { OwnedFd::from_raw_fd(fd) });
Bell { fd }
}
}
impl Bell {
pub(super) fn ring(&self) {
if let Some(fd) = &self.fd {
let one = 1u64.to_ne_bytes();
unsafe { libc::write(fd.as_raw_fd(), one.as_ptr().cast(), one.len()) };
}
}
fn hush(&self) {
if let Some(fd) = &self.fd {
let mut count = [0u8; 8];
unsafe { libc::read(fd.as_raw_fd(), count.as_mut_ptr().cast(), count.len()) };
}
}
}
const REMOTE_FILE_SYSTEMS: [i64; 11] = [
0x6969, 0x517B, 0xFF53_4D42_u32 as i64, 0xFE53_4D42_u32 as i64, 0x6573_5546, 0x00C3_6400, 0x5346_414F, 0x6B41_4653, 0x0102_1997, 0x7375_7245, 0x0BD0_0BD0, ];
fn remote(path: &Path) -> bool {
let Ok(c_path) = std::ffi::CString::new(path.as_os_str().as_encoded_bytes()) else {
return true;
};
let mut stat: libc::statfs = unsafe { std::mem::zeroed() };
if unsafe { libc::statfs(c_path.as_ptr(), &mut stat) } != 0 {
return true;
}
let kind = stat.f_type as i64 & 0xFFFF_FFFF;
REMOTE_FILE_SYSTEMS.contains(&kind)
}
pub(super) struct Notify {
fd: OwnedFd,
}
const FILE_EVENTS: u32 = libc::IN_MODIFY
| libc::IN_ATTRIB
| libc::IN_CLOSE_WRITE
| libc::IN_MOVE_SELF
| libc::IN_DELETE_SELF;
const DIRECTORY_EVENTS: u32 =
libc::IN_CREATE | libc::IN_MOVED_TO | libc::IN_MOVED_FROM | libc::IN_DELETE | libc::IN_MODIFY;
impl Notify {
pub(super) fn new(path: &Path) -> Option<Notify> {
if remote(path) {
return None;
}
let fd = unsafe { libc::inotify_init1(libc::IN_NONBLOCK | libc::IN_CLOEXEC) };
if fd < 0 {
return None;
}
let notify = Notify {
fd: unsafe { OwnedFd::from_raw_fd(fd) },
};
let directory = match path.parent() {
Some(parent) if !parent.as_os_str().is_empty() => parent,
_ => Path::new("."),
};
(notify.watch(path, FILE_EVENTS) && notify.watch(directory, DIRECTORY_EVENTS))
.then_some(notify)
}
pub(super) fn rewatch(&self, path: &Path) {
self.watch(path, FILE_EVENTS);
}
fn watch(&self, path: &Path, mask: u32) -> bool {
let Ok(c_path) = std::ffi::CString::new(path.as_os_str().as_encoded_bytes()) else {
return false;
};
unsafe { libc::inotify_add_watch(self.fd.as_raw_fd(), c_path.as_ptr(), mask) >= 0 }
}
pub(super) fn wait(&self, bell: &Bell, timeout: Option<std::time::Duration>) -> Woke {
let mut fds = [
libc::pollfd {
fd: self.fd.as_raw_fd(),
events: libc::POLLIN,
revents: 0,
},
libc::pollfd {
fd: bell.fd.as_ref().map_or(-1, |fd| fd.as_raw_fd()),
events: libc::POLLIN,
revents: 0,
},
];
let timeout = timeout.map_or(-1, |t| i32::try_from(t.as_millis()).unwrap_or(i32::MAX));
loop {
let n = unsafe { libc::poll(fds.as_mut_ptr(), fds.len() as libc::nfds_t, timeout) };
if n < 0 && std::io::Error::last_os_error().kind() == std::io::ErrorKind::Interrupted {
continue;
}
if n == 0 {
return Woke::TimedOut;
}
if fds[1].revents != 0 {
bell.hush();
return Woke::Rung;
}
return Woke::Changed;
}
}
pub(super) fn drain(&self) {
let mut buf = [0u8; 4096];
while unsafe { libc::read(self.fd.as_raw_fd(), buf.as_mut_ptr().cast(), buf.len()) } > 0 {}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write as _;
use std::time::Duration;
#[test]
fn a_wait_wakes_for_changes_and_the_bell() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("log.csv");
std::fs::write(&path, "t\n1\n").unwrap();
let notify = Notify::new(&path).expect("tmp takes an inotify watch");
let bell = Bell::default();
let short = Some(Duration::from_millis(20));
assert_eq!(notify.wait(&bell, short), Woke::TimedOut);
let mut file = std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap();
file.write_all(b"2\n").unwrap();
let guard = Some(Duration::from_secs(30));
assert_eq!(notify.wait(&bell, guard), Woke::Changed);
notify.drain();
assert_eq!(notify.wait(&bell, short), Woke::TimedOut, "drained");
let next = dir.path().join("next.csv");
std::fs::write(&next, "t\n9\n").unwrap();
notify.drain();
std::fs::rename(&next, &path).unwrap();
assert_eq!(notify.wait(&bell, guard), Woke::Changed);
notify.drain();
bell.ring();
assert_eq!(notify.wait(&bell, guard), Woke::Rung);
assert_eq!(notify.wait(&bell, short), Woke::TimedOut, "hushed");
}
}