use interprocess::local_socket::{
GenericFilePath, ListenerNonblockingMode, ListenerOptions, ToFsName,
};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::fmt;
use std::io;
use std::path::{Path, PathBuf};
#[derive(Debug)]
pub enum LayoutError {
Io { path: PathBuf, source: io::Error },
}
impl fmt::Display for LayoutError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Io { path, source } => write!(f, "{}: {source}", path.display()),
}
}
}
impl std::error::Error for LayoutError {}
pub fn apply_private_mode(path: &Path) -> Result<(), LayoutError> {
set_file_mode(path, 0o600).map_err(|source| LayoutError::Io {
path: path.to_path_buf(),
source,
})
}
pub const UNIX_SOCKET_PATH_MAX: usize = 103;
pub const SOCKET_FILE_NAME: &str = "s";
pub const SOCKET_SUFFIX: &str = ".sock";
pub const REGISTRATION_SUFFIX: &str = ".json";
pub const RUNTIME_DIR_ENV: &str = "ONLYNE_RUNTIME_DIR";
pub const WIRE_VERSION: &str = env!("CARGO_PKG_VERSION");
pub fn runtime_dir() -> io::Result<PathBuf> {
let dir = runtime_dir_path();
ensure_private_dir(&dir)?;
Ok(dir)
}
pub fn runtime_dir_path() -> PathBuf {
match std::env::var_os(RUNTIME_DIR_ENV) {
Some(dir) if !dir.is_empty() => PathBuf::from(dir),
_ => runtime_base().join(format!("onlyne-{}", current_user_id())),
}
}
#[cfg(unix)]
fn runtime_base() -> PathBuf {
PathBuf::from("/tmp")
}
#[cfg(not(unix))]
fn runtime_base() -> PathBuf {
std::env::temp_dir()
}
#[cfg(unix)]
fn current_user_id() -> String {
format!("{}", unsafe { libc::geteuid() })
}
#[cfg(not(unix))]
fn current_user_id() -> String {
"user".to_string()
}
pub fn socket_path(root: &Path) -> io::Result<PathBuf> {
runtime_dir()?;
Ok(runtime_file_path(root, SOCKET_SUFFIX))
}
pub fn registration_path(root: &Path) -> PathBuf {
runtime_file_path(root, REGISTRATION_SUFFIX)
}
fn runtime_file_path(root: &Path, suffix: &str) -> PathBuf {
let digest = workspace_digest(root);
runtime_dir_path().join(format!("{digest}{suffix}"))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RegistrationKind {
Server,
Client,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RegistrationFile {
pub kind: RegistrationKind,
#[serde(default)]
pub role: Option<String>,
pub root: PathBuf,
pub pid: u32,
#[serde(default)]
pub version: String,
#[serde(default)]
pub runtime: Option<String>,
#[serde(default)]
pub placement: Option<String>,
}
impl RegistrationFile {
pub fn server(root: &Path) -> Self {
Self::new(RegistrationKind::Server, root)
}
pub fn client(root: &Path) -> Self {
Self::new(RegistrationKind::Client, root)
}
pub fn new(kind: RegistrationKind, root: &Path) -> Self {
Self {
kind,
role: None,
root: absolute_path(root),
pid: std::process::id(),
version: WIRE_VERSION.to_string(),
runtime: None,
placement: None,
}
}
pub fn with_role(mut self, role: impl Into<String>) -> Self {
self.role = Some(role.into());
self
}
pub fn with_runtime(mut self, runtime: impl Into<String>) -> Self {
self.runtime = Some(runtime.into());
self
}
pub fn with_placement(mut self, placement: impl Into<String>) -> Self {
self.placement = Some(placement.into());
self
}
}
pub fn write_registration(root: &Path, reg: &RegistrationFile) -> io::Result<()> {
let path = runtime_dir()?.join(format!("{}{REGISTRATION_SUFFIX}", workspace_digest(root)));
write_registration_at(&path, reg)
}
pub fn read_registration(root: &Path) -> io::Result<Option<RegistrationFile>> {
read_registration_at(®istration_path(root))
}
pub fn remove_registration(root: &Path) -> io::Result<()> {
match std::fs::remove_file(registration_path(root)) {
Ok(()) => Ok(()),
Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(error),
}
}
pub fn list_registrations() -> io::Result<Vec<(PathBuf, RegistrationFile)>> {
let entries = match std::fs::read_dir(runtime_dir_path()) {
Ok(entries) => entries,
Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(error) => return Err(error),
};
let mut found = Vec::new();
for entry in entries {
let path = entry?.path();
if !path
.to_str()
.is_some_and(|name| name.ends_with(REGISTRATION_SUFFIX))
{
continue;
}
if let Ok(Some(reg)) = read_registration_at(&path) {
found.push((path, reg));
}
}
found.sort_by(|left, right| left.0.cmp(&right.0));
Ok(found)
}
fn write_registration_at(path: &Path, reg: &RegistrationFile) -> io::Result<()> {
let json = serde_json::to_vec_pretty(reg)
.map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
write_private_file(path, &json)
}
fn read_registration_at(path: &Path) -> io::Result<Option<RegistrationFile>> {
match std::fs::read(path) {
Ok(bytes) => match serde_json::from_slice(&bytes) {
Ok(reg) => Ok(Some(reg)),
Err(error) => Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("{}: {error}", path.display()),
)),
},
Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(None),
Err(error) => Err(error),
}
}
fn write_private_file(path: &Path, bytes: &[u8]) -> io::Result<()> {
use std::io::Write as _;
let name = path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("onlyne-registration");
let temp = path.with_file_name(format!("{name}.tmp-{}", std::process::id()));
let outcome = open_private(&temp).and_then(|mut file| {
file.write_all(bytes)?;
drop(file);
std::fs::rename(&temp, path)
});
if outcome.is_err() {
let _ = std::fs::remove_file(&temp);
}
outcome
}
#[cfg(unix)]
fn open_private(path: &Path) -> io::Result<std::fs::File> {
use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
let file = std::fs::OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.mode(0o600)
.open(path)?;
let mut permissions = file.metadata()?.permissions();
permissions.set_mode(0o600);
file.set_permissions(permissions)?;
Ok(file)
}
#[cfg(not(unix))]
fn open_private(path: &Path) -> io::Result<std::fs::File> {
std::fs::OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.open(path)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SocketEndpoint {
root: PathBuf,
natural: PathBuf,
actual: PathBuf,
registration: PathBuf,
}
impl SocketEndpoint {
pub fn resolve(root: &Path, run_dir: &Path) -> Self {
let root = absolute_path(root);
let dir = runtime_dir_path();
let digest = workspace_digest(&root);
Self {
natural: run_dir.join(SOCKET_FILE_NAME),
actual: dir.join(format!("{digest}{SOCKET_SUFFIX}")),
registration: dir.join(format!("{digest}{REGISTRATION_SUFFIX}")),
root,
}
}
pub fn root(&self) -> &Path {
&self.root
}
pub fn natural(&self) -> &Path {
&self.natural
}
pub fn actual(&self) -> &Path {
&self.actual
}
pub fn registration(&self) -> &Path {
&self.registration
}
pub fn short(&self) -> bool {
self.actual != self.natural
}
pub fn publish(&self, reg: &RegistrationFile) -> io::Result<()> {
runtime_dir()?;
write_registration_at(&self.registration, reg)
}
}
pub fn absolute_path(path: &Path) -> PathBuf {
if let Ok(canonical) = std::fs::canonicalize(path) {
return canonical;
}
collapse_parents(&lexical_absolute(path))
}
fn collapse_parents(path: &Path) -> PathBuf {
let mut out = PathBuf::new();
for component in path.components() {
match component {
std::path::Component::ParentDir => match out.components().next_back() {
Some(std::path::Component::Normal(_)) => {
out.pop();
}
_ => out.push(component.as_os_str()),
},
other => out.push(other.as_os_str()),
}
}
out
}
pub fn create_dir(path: &Path, mode: Option<u32>) -> io::Result<()> {
std::fs::create_dir_all(path)?;
if let Some(mode) = mode {
set_dir_mode(path, mode)?;
}
Ok(())
}
#[cfg(unix)]
pub fn set_dir_mode(path: &Path, mode: u32) -> io::Result<()> {
use std::os::unix::fs::PermissionsExt;
let mut permissions = std::fs::metadata(path)?.permissions();
permissions.set_mode(mode);
std::fs::set_permissions(path, permissions)
}
#[cfg(not(unix))]
pub fn set_dir_mode(_path: &Path, _mode: u32) -> io::Result<()> {
Ok(())
}
fn set_file_mode(path: &Path, mode: u32) -> io::Result<()> {
let _ = std::fs::metadata(path)?;
set_file_mode_impl(path, mode)
}
#[cfg(unix)]
fn set_file_mode_impl(path: &Path, mode: u32) -> io::Result<()> {
use std::os::unix::fs::PermissionsExt;
let mut permissions = std::fs::metadata(path)?.permissions();
permissions.set_mode(mode);
std::fs::set_permissions(path, permissions)
}
#[cfg(not(unix))]
fn set_file_mode_impl(_path: &Path, _mode: u32) -> io::Result<()> {
Ok(())
}
pub type LocalListener = interprocess::local_socket::tokio::Listener;
pub type LocalStream = interprocess::local_socket::tokio::Stream;
pub type LocalListenerSync = interprocess::local_socket::Listener;
pub type LocalStreamSync = interprocess::local_socket::Stream;
pub mod prelude {
pub use interprocess::local_socket::traits::tokio::{
Listener as TokioListener, Stream as TokioStream,
};
pub use interprocess::local_socket::traits::{Listener as SyncListener, Stream as SyncStream};
}
#[cfg(windows)]
const MARKER_PREFIX: &str = "v1:";
const PIPE_BUSY: i32 = 231;
const VERBATIM_PIPE_PREFIX: &str = r"\\.\pipe\";
pub fn is_verbatim_pipe_path(path: &Path) -> bool {
path.to_str()
.is_some_and(|s| starts_with_ignore_ascii_case(s, VERBATIM_PIPE_PREFIX))
}
fn lexical_absolute(path: &Path) -> PathBuf {
std::path::absolute(path).unwrap_or_else(|_| path.to_path_buf())
}
fn digest_hex(path: &Path, bytes: usize) -> String {
let normalized = path
.to_string_lossy()
.replace('\\', "/")
.to_ascii_lowercase();
let digest = Sha256::digest(normalized.as_bytes());
let mut hex = String::with_capacity(bytes * 2);
for byte in &digest[..bytes] {
hex.push(HEX[(*byte >> 4) as usize] as char);
hex.push(HEX[(*byte & 0x0f) as usize] as char);
}
hex
}
pub fn pipe_name_for(path: &Path) -> String {
format!("onlyne-{}", digest_hex(&lexical_absolute(path), 16))
}
pub fn workspace_digest(root: &Path) -> String {
digest_hex(&absolute_path(root), 8)
}
const HEX: &[u8; 16] = b"0123456789abcdef";
#[allow(clippy::unused_async)]
pub async fn bind_local(path: &Path) -> io::Result<LocalListener> {
bind_tokio(path)
}
pub fn bind_tokio(path: &Path) -> io::Result<LocalListener> {
create_with_privacy(
path,
ListenerNonblockingMode::Neither,
ListenerOptions::create_tokio,
)
}
pub fn bind_socket_v2(root: &Path) -> io::Result<LocalListener> {
let path = socket_path(root)?;
bind_runtime_socket(&path)
}
pub fn bind_socket(root: &Path, run_dir: &Path) -> io::Result<(LocalListener, SocketEndpoint)> {
let endpoint = SocketEndpoint::resolve(root, run_dir);
runtime_dir()?;
create_dir(run_dir, Some(0o700))?;
let listener = bind_runtime_socket(endpoint.actual())?;
Ok((listener, endpoint))
}
pub fn bind_socket_registered(
root: &Path,
run_dir: &Path,
reg: &RegistrationFile,
) -> io::Result<(LocalListener, SocketEndpoint)> {
let (listener, endpoint) = bind_socket(root, run_dir)?;
endpoint.publish(reg)?;
Ok((listener, endpoint))
}
fn bind_runtime_socket(path: &Path) -> io::Result<LocalListener> {
#[cfg(unix)]
remove_stale_socket(path);
bind_tokio(path).map_err(|error| bind_failure(path, error))
}
#[cfg(unix)]
fn remove_stale_socket(path: &Path) {
let _ = std::fs::remove_file(path);
}
fn bind_failure(path: &Path, error: io::Error) -> io::Error {
let dir = path.parent().unwrap_or(Path::new(""));
io::Error::new(
error.kind(),
format!(
"bind {} ({} bytes) in {}: {error}",
path.display(),
path.as_os_str().len(),
dir.display(),
),
)
}
#[cfg(unix)]
fn ensure_private_dir(path: &Path) -> io::Result<()> {
match std::fs::symlink_metadata(path) {
Ok(metadata) => verify_private_dir(path, &metadata)?,
Err(error) if error.kind() == io::ErrorKind::NotFound => {
std::fs::create_dir_all(path).map_err(|source| dir_failure(path, &source))?;
set_dir_mode(path, 0o700).map_err(|source| dir_failure(path, &source))?;
}
Err(source) => return Err(dir_failure(path, &source)),
}
static PROBE_SEQUENCE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let sequence = PROBE_SEQUENCE.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let probe = path.join(format!("onlyne-probe-{}-{sequence}", std::process::id()));
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&probe)
{
Ok(_) => {
let _ = std::fs::remove_file(&probe);
Ok(())
}
Err(source) => Err(dir_failure(path, &source)),
}
}
#[cfg(unix)]
fn verify_private_dir(path: &Path, metadata: &std::fs::Metadata) -> io::Result<()> {
use std::os::unix::fs::{MetadataExt, PermissionsExt};
if metadata.file_type().is_symlink() {
return Err(dir_refusal(path, "a symlink"));
}
if !metadata.is_dir() {
return Err(dir_refusal(path, "already held by a non-directory"));
}
if metadata.uid() != unsafe { libc::geteuid() } {
return Err(dir_refusal(path, "owned by another user"));
}
if metadata.permissions().mode() & 0o077 != 0 {
set_dir_mode(path, 0o700).map_err(|source| dir_failure(path, &source))?;
}
Ok(())
}
#[cfg(not(unix))]
fn ensure_private_dir(path: &Path) -> io::Result<()> {
if path.is_dir() {
return Ok(());
}
std::fs::create_dir_all(path)
}
#[cfg(unix)]
fn dir_refusal(path: &Path, reason: &str) -> io::Error {
io::Error::new(
io::ErrorKind::PermissionDenied,
format!("refusing to use directory {}: {reason}", path.display()),
)
}
#[cfg(unix)]
fn dir_failure(path: &Path, source: &io::Error) -> io::Error {
io::Error::new(
source.kind(),
format!("directory {}: {source}", path.display()),
)
}
pub fn bind_local_sync(path: &Path) -> io::Result<LocalListenerSync> {
create_with_privacy(
path,
ListenerNonblockingMode::Neither,
ListenerOptions::create_sync,
)
}
pub fn bind_local_sync_poll(path: &Path) -> io::Result<LocalListenerSync> {
create_with_privacy(
path,
ListenerNonblockingMode::Accept,
ListenerOptions::create_sync,
)
}
pub async fn connect_local(path: &Path) -> io::Result<LocalStream> {
use interprocess::local_socket::tokio::Stream;
use interprocess::local_socket::traits::tokio::Stream as _;
Stream::connect(connect_name(path)?)
.await
.map_err(map_pipe_busy)
}
pub fn connect_local_sync(path: &Path) -> io::Result<LocalStreamSync> {
use interprocess::local_socket::Stream;
use interprocess::local_socket::traits::Stream as _;
Stream::connect(connect_name(path)?).map_err(map_pipe_busy)
}
fn map_pipe_busy(err: io::Error) -> io::Error {
if err.raw_os_error() == Some(PIPE_BUSY) {
io::Error::new(io::ErrorKind::WouldBlock, err)
} else {
err
}
}
fn starts_with_ignore_ascii_case(value: &str, prefix: &str) -> bool {
value.len() >= prefix.len()
&& value
.as_bytes()
.iter()
.zip(prefix.as_bytes())
.all(|(a, b)| a.eq_ignore_ascii_case(b))
}
#[cfg(windows)]
fn parse_marker(text: &str) -> Option<String> {
let rest = text.trim().strip_prefix(MARKER_PREFIX)?;
if rest.is_empty()
|| !rest
.bytes()
.all(|b| b.is_ascii_graphic() && b != b'/' && b != b'\\')
{
return None;
}
Some(rest.to_string())
}
#[cfg(windows)]
fn read_marker_name(path: &Path) -> Option<String> {
let text = std::fs::read_to_string(path).ok()?;
parse_marker(&text)
}
#[cfg(windows)]
fn write_marker(path: &Path, pipe_name: &str) -> io::Result<()> {
std::fs::write(path, format!("{MARKER_PREFIX}{pipe_name}"))
}
#[cfg(windows)]
fn windows_pipe_leaf(path: &Path) -> String {
read_marker_name(path).unwrap_or_else(|| pipe_name_for(path))
}
fn create_with_privacy<T>(
path: &Path,
nonblocking: ListenerNonblockingMode,
create: fn(ListenerOptions<'static>) -> io::Result<T>,
) -> io::Result<T> {
#[cfg(unix)]
{
use interprocess::os::unix::local_socket::ListenerOptionsExt;
match create(unix_options(path, nonblocking)?.mode(0o600)) {
Ok(listener) => Ok(listener),
Err(err) if err.kind() == io::ErrorKind::Unsupported => {
let listener = create(unix_options(path, nonblocking)?)?;
apply_private_mode(path).map_err(|e| io::Error::other(e.to_string()))?;
Ok(listener)
}
Err(err) => Err(err),
}
}
#[cfg(windows)]
{
create(listener_options_windows(path, nonblocking)?)
}
}
#[cfg(unix)]
fn unix_options(
path: &Path,
nonblocking: ListenerNonblockingMode,
) -> io::Result<ListenerOptions<'static>> {
let name = path.to_fs_name::<GenericFilePath>()?.into_owned();
Ok(ListenerOptions::new()
.name(name)
.reclaim_name(false)
.nonblocking(nonblocking))
}
#[cfg(windows)]
fn listener_options_windows(
path: &Path,
nonblocking: ListenerNonblockingMode,
) -> io::Result<ListenerOptions<'static>> {
use interprocess::local_socket::{GenericNamespaced, ToNsName};
use interprocess::os::windows::local_socket::ListenerOptionsExt;
use interprocess::os::windows::security_descriptor::SecurityDescriptor;
use widestring::u16cstr;
let name = if is_verbatim_pipe_path(path) {
path.to_fs_name::<GenericFilePath>()?.into_owned()
} else {
let leaf = windows_pipe_leaf(path);
write_marker(path, &leaf)?;
leaf.to_ns_name::<GenericNamespaced>()?.into_owned()
};
let sd = SecurityDescriptor::deserialize(u16cstr!("D:P(A;;GA;;;OW)(A;;GA;;;SY)"))?;
Ok(ListenerOptions::new()
.name(name)
.reclaim_name(false)
.nonblocking(nonblocking)
.security_descriptor(sd))
}
fn connect_name(path: &Path) -> io::Result<interprocess::local_socket::Name<'static>> {
#[cfg(unix)]
{
Ok(path.to_fs_name::<GenericFilePath>()?.into_owned())
}
#[cfg(windows)]
{
use interprocess::local_socket::{GenericNamespaced, ToNsName};
if is_verbatim_pipe_path(path) {
Ok(path.to_fs_name::<GenericFilePath>()?.into_owned())
} else {
Ok(windows_pipe_leaf(path)
.to_ns_name::<GenericNamespaced>()?
.into_owned())
}
}
}