Skip to main content

onlyne_wire/
socket.rs

1//! Cross-platform local socket seam, the owner-tree endpoint over it, and the
2//! per-user runtime directory that names both.
3//!
4//! v2 keeps every local socket in one machine-level runtime directory:
5//! `$ONLYNE_RUNTIME_DIR` when the operator sets it, `/tmp/onlyne-<uid>/`
6//! otherwise, created `0700` and adopted only when this process owns it alone.
7//! One workspace owns two files there, both named by [`workspace_digest`] — the
8//! first 16 hex characters of `sha256` over the canonical workspace root:
9//!
10//! - `<digest>.sock`: the socket [`bind_socket_v2`] binds and [`connect_local`]
11//!   reaches, `0600`.
12//! - `<digest>.json`: the [`RegistrationFile`] naming who serves it — kind,
13//!   role, root, pid, version, runtime. [`read_registration`],
14//!   [`registration_path`], and [`list_registrations`] are the discovery seam
15//!   the CLI, the test harness, and external-runtime plugins read instead of
16//!   walking workspace trees.
17//!
18//! v1 kept the socket at `<run_dir>/s` while that spelling fit `sun_path`, moved
19//! it to a short derived path when it did not, and published the choice in
20//! `run/socket`. A generated role nests deep enough
21//! (`<root>/.onlyne/ws/<topology>/<role>/.onlyne/run/s`) that the canonical
22//! spelling passes [`UNIX_SOCKET_PATH_MAX`], so the two-rule answer split one
23//! tree across two directories depending on who resolved when. The runtime path
24//! is about 40 bytes for any root, so one tree owns one spelling, and
25//! `<run_dir>/s` survives only as [`SocketEndpoint::natural`], the spelling
26//! operators and older tooling print.
27//!
28//! The unix base is `/tmp` rather than [`std::env::temp_dir`] because a
29//! launchd-started daemon and an interactive shell see different `TMPDIR` and
30//! `XDG_RUNTIME_DIR` values; one fixed base makes every context compute the same
31//! path for the same root.
32//!
33//! Windows cannot bind a UDS on stable 1.85, so `<digest>.sock` there is a
34//! regular marker file whose contents name an NPFS pipe (`v1:onlyne-<32hex>`)
35//! and the pipe carries the traffic. `--socket` values that already start with a
36//! verbatim pipe path skip derivation and travel as `GenericFilePath`. Frame I/O
37//! stays generic `AsyncRead`/`AsyncWrite`.
38//!
39//! [`create_dir`], [`set_dir_mode`], and [`apply_private_mode`] are the mode half
40//! of the same seam: one spelling of "owner-only `run/`, `keys/`, and runtime
41//! directory" for the daemon that binds and the layout that bootstraps.
42
43use interprocess::local_socket::{
44    GenericFilePath, ListenerNonblockingMode, ListenerOptions, ToFsName,
45};
46use serde::{Deserialize, Serialize};
47use sha2::{Digest, Sha256};
48use std::fmt;
49use std::io;
50use std::path::{Path, PathBuf};
51
52/// Errors produced by the owner-only mode helpers below.
53#[derive(Debug)]
54pub enum LayoutError {
55    /// Filesystem operation failed for a concrete path.
56    Io { path: PathBuf, source: io::Error },
57}
58
59impl fmt::Display for LayoutError {
60    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
61        match self {
62            Self::Io { path, source } => write!(f, "{}: {source}", path.display()),
63        }
64    }
65}
66
67impl std::error::Error for LayoutError {}
68
69/// Apply `0600` to an existing private file.
70///
71/// Bootstrap creates `run/` and `keys/` with `0700`. Daemons call this helper
72/// after writing key files. Socket privacy is applied at bind time by
73/// [`bind_local`] (unix `mode(0o600)`, windows owner-only SDDL), and
74/// registration privacy by [`write_registration`], which removes the chmod
75/// TOCTOU a post-write call here would have.
76pub fn apply_private_mode(path: &Path) -> Result<(), LayoutError> {
77    set_file_mode(path, 0o600).map_err(|source| LayoutError::Io {
78        path: path.to_path_buf(),
79        source,
80    })
81}
82
83/// Bytes available in `sun_path` on unix, including the trailing NUL. macOS
84/// allows 104.
85///
86/// Every v2 runtime path is about 40 bytes, so no caller has to respect this
87/// bound any more; the constant stays because it is what v1 worked around. A
88/// canonical spelling over it is not a legal name for the kernel: the bind fails
89/// and `connect()` fails the same way for every client that keeps retrying the
90/// spelling it derived — the reported v1 symptoms were a client logging
91/// `adapter socket restarting error=bind <path>` on a half-second loop while
92/// `onlyne status` kept reporting the role as connected, because the TLS link to
93/// the server was healthy and the local half was dead. A generated role
94/// workspace nests three levels below its server root, so a root that was
95/// already long carried the canonical socket spelling past 103 bytes.
96pub const UNIX_SOCKET_PATH_MAX: usize = 103;
97
98/// Socket file name inside `run/`, the canonical spelling every v1 layout path
99/// used. v2 binds elsewhere and keeps this only as
100/// [`SocketEndpoint::natural`].
101pub const SOCKET_FILE_NAME: &str = "s";
102
103/// Socket leaf inside a runtime directory: `<digest>.sock`.
104pub const SOCKET_SUFFIX: &str = ".sock";
105
106/// Registration leaf inside a runtime directory: `<digest>.json`.
107pub const REGISTRATION_SUFFIX: &str = ".json";
108
109/// Environment variable that replaces the default runtime directory, used
110/// verbatim when set and non-empty.
111pub const RUNTIME_DIR_ENV: &str = "ONLYNE_RUNTIME_DIR";
112
113/// The version of the wire protocol this crate speaks, the `version` every
114/// [`RegistrationFile`] written here carries, so a reader can tell which
115/// protocol a live endpoint expects.
116pub const WIRE_VERSION: &str = env!("CARGO_PKG_VERSION");
117
118/// The runtime directory, created and owner-verified.
119///
120/// `$ONLYNE_RUNTIME_DIR` when set and non-empty, `/tmp/onlyne-<uid>/`
121/// otherwise. Adoption follows [`ensure_private_dir`]: a directory this user
122/// owns that is closed to group and other access, which this process can still
123/// write. `/tmp` is shared and world-executable, so another user can create
124/// `/tmp/onlyne-<uid>` first; an endpoint inside a directory that user owns is
125/// one that user could swap from under a live daemon, which is the refusal
126/// [`ensure_private_dir`] makes.
127pub fn runtime_dir() -> io::Result<PathBuf> {
128    let dir = runtime_dir_path();
129    ensure_private_dir(&dir)?;
130    Ok(dir)
131}
132
133/// The runtime directory's spelling without creating or checking anything, so a
134/// caller can name a path before any daemon exists (`/tmp` is one component).
135pub fn runtime_dir_path() -> PathBuf {
136    match std::env::var_os(RUNTIME_DIR_ENV) {
137        Some(dir) if !dir.is_empty() => PathBuf::from(dir),
138        _ => runtime_base().join(format!("onlyne-{}", current_user_id())),
139    }
140}
141
142/// `/tmp` on unix: fixed so launchd and an interactive shell agree, and short
143/// enough that every derived path stays well under [`UNIX_SOCKET_PATH_MAX`].
144#[cfg(unix)]
145fn runtime_base() -> PathBuf {
146    PathBuf::from("/tmp")
147}
148
149/// Off unix the per-user temporary root already scopes the directory.
150#[cfg(not(unix))]
151fn runtime_base() -> PathBuf {
152    std::env::temp_dir()
153}
154
155/// The effective uid, the runtime directory's name component.
156#[cfg(unix)]
157fn current_user_id() -> String {
158    // SAFETY: `geteuid` takes no arguments, touches no memory, and cannot fail.
159    format!("{}", unsafe { libc::geteuid() })
160}
161
162/// Off unix the per-user temporary root is the scope: `%TEMP%` sits inside the
163/// user profile, and the pipe's owner-only security descriptor guards the
164/// endpoint, so no uid is spelled out.
165#[cfg(not(unix))]
166fn current_user_id() -> String {
167    "user".to_string()
168}
169
170/// The socket one workspace owns: `<runtime_dir>/<workspace_digest>.sock`.
171///
172/// The runtime directory is created and verified first, so a caller that gets a
173/// path has a private place to bind it in.
174pub fn socket_path(root: &Path) -> io::Result<PathBuf> {
175    runtime_dir()?;
176    Ok(runtime_file_path(root, SOCKET_SUFFIX))
177}
178
179/// The registration belonging to [`socket_path`]'s socket,
180/// `<runtime_dir>/<workspace_digest>.json`.
181///
182/// Pure: [`write_registration`] and [`socket_path`] create the directory, so a
183/// reader can name the file without claiming the machine has a runtime.
184pub fn registration_path(root: &Path) -> PathBuf {
185    runtime_file_path(root, REGISTRATION_SUFFIX)
186}
187
188/// `<runtime_dir_path>/<workspace_digest(root)><suffix>`, filesystem untouched.
189fn runtime_file_path(root: &Path, suffix: &str) -> PathBuf {
190    let digest = workspace_digest(root);
191    runtime_dir_path().join(format!("{digest}{suffix}"))
192}
193
194/// Which side of a socket one [`RegistrationFile`] describes.
195///
196/// The serving side owns the answer: a server root's daemon publishes
197/// [`Server`](Self::Server), and a role workspace's client daemon publishes
198/// [`Client`](Self::Client) with its role. No reader infers the surface from
199/// what happens to sit beside a tree — two v1 readers did, and disagreed.
200#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
201#[serde(rename_all = "snake_case")]
202pub enum RegistrationKind {
203    /// The admin surface of a server root.
204    Server,
205    /// A role workspace's client daemon.
206    Client,
207}
208
209/// The one file that names a workspace's endpoint: `<runtime>/<digest>.json`.
210///
211/// It answers without touching the workspace what `run/socket` answered for a
212/// path and what a `state.db`/`client.db` guess answered for a surface: who
213/// serves this tree (`kind`, `role`), which tree it is (`root`), which process
214/// (`pid`), which wire version, and which runtime hosts the role's sessions
215/// (`runtime`, the field external-runtime plugins match to find their clients).
216///
217/// The writer's facts are the file's facts. A bind does not publish one on its
218/// own, because a bind does not know whether it belongs to an admin root or a
219/// role workspace: the serving side calls [`write_registration`] or
220/// [`SocketEndpoint::publish`] with the surface it knows.
221#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
222pub struct RegistrationFile {
223    /// Which side of the socket this process serves.
224    pub kind: RegistrationKind,
225    /// The role this client serves; `None` for a server root. `null` in JSON.
226    #[serde(default)]
227    pub role: Option<String>,
228    /// The canonical spelling of the tree the socket belongs to.
229    pub root: PathBuf,
230    /// The pid of the process serving it, restated at every bind.
231    pub pid: u32,
232    /// The writer's version, [`WIRE_VERSION`] for this crate.
233    #[serde(default)]
234    pub version: String,
235    /// The runtime hosting the role's sessions (`pi`, `acp`, …); `None` for a
236    /// server root. `null` in JSON.
237    #[serde(default)]
238    pub runtime: Option<String>,
239    /// Where this machine displays the role's runtime process (`orca`, `zellij`,
240    /// `headless`, `external`); `None` for a server root, and for a client that
241    /// published its registration before it resolved one.
242    ///
243    /// Placement is a property of the machine, which is why it is recorded
244    /// here: an external runtime's plugin reads the registrations in this
245    /// directory and dials the clients whose placement says the runtime is
246    /// already resident (`external`).
247    #[serde(default)]
248    pub placement: Option<String>,
249}
250
251impl RegistrationFile {
252    /// This process serving `root` as a server root.
253    pub fn server(root: &Path) -> Self {
254        Self::new(RegistrationKind::Server, root)
255    }
256
257    /// This process serving `root` as a role workspace's client daemon.
258    pub fn client(root: &Path) -> Self {
259        Self::new(RegistrationKind::Client, root)
260    }
261
262    /// This process, `WIRE_VERSION`, canonical `root`, no role, no runtime.
263    pub fn new(kind: RegistrationKind, root: &Path) -> Self {
264        Self {
265            kind,
266            role: None,
267            root: absolute_path(root),
268            pid: std::process::id(),
269            version: WIRE_VERSION.to_string(),
270            runtime: None,
271            placement: None,
272        }
273    }
274
275    /// Name the role this registration serves.
276    pub fn with_role(mut self, role: impl Into<String>) -> Self {
277        self.role = Some(role.into());
278        self
279    }
280
281    /// Name the runtime hosting the role's sessions.
282    pub fn with_runtime(mut self, runtime: impl Into<String>) -> Self {
283        self.runtime = Some(runtime.into());
284        self
285    }
286
287    /// Name the placement this machine displays the role's sessions in.
288    pub fn with_placement(mut self, placement: impl Into<String>) -> Self {
289        self.placement = Some(placement.into());
290        self
291    }
292}
293
294/// Publish `reg` as the registration for `root`, creating the runtime directory.
295///
296/// The digest is taken now, so a caller that has just bound should publish
297/// through [`SocketEndpoint::publish`] or bind with [`bind_socket_registered`],
298/// which reuse the digest the bind resolved.
299///
300/// The file is replaced atomically, so a reader listing the directory never sees
301/// half a registration.
302pub fn write_registration(root: &Path, reg: &RegistrationFile) -> io::Result<()> {
303    let path = runtime_dir()?.join(format!("{}{REGISTRATION_SUFFIX}", workspace_digest(root)));
304    write_registration_at(&path, reg)
305}
306
307/// The registration for `root`, or `None` when nothing published one.
308///
309/// Only absence means "nothing is published". A file that exists at this tree's
310/// own digest but is not a registration is [`io::ErrorKind::InvalidData`] with
311/// the path in the message, because a reader that got this far has to know it
312/// could not be read rather than treat a broken endpoint as an absent one.
313pub fn read_registration(root: &Path) -> io::Result<Option<RegistrationFile>> {
314    read_registration_at(&registration_path(root))
315}
316
317/// Remove the registration for `root`. Removing nothing is not an error.
318pub fn remove_registration(root: &Path) -> io::Result<()> {
319    match std::fs::remove_file(registration_path(root)) {
320        Ok(()) => Ok(()),
321        Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
322        Err(error) => Err(error),
323    }
324}
325
326/// Every registration this machine's runtime directory holds, sorted by file
327/// path.
328///
329/// A missing runtime directory is no registrations rather than an error: a
330/// reader that starts before any daemon should see an empty machine. Files that
331/// are not `.json`, and files that do not parse as a [`RegistrationFile`], are
332/// skipped — a publisher writes through a temporary name and holds the socket in
333/// a private directory, and one stray file must not blind a reader to every
334/// other tree. The socket each entry describes is `socket_path(&reg.root)`, or
335/// the entry's own path with its `.json` leaf replaced by `.sock`.
336pub fn list_registrations() -> io::Result<Vec<(PathBuf, RegistrationFile)>> {
337    let entries = match std::fs::read_dir(runtime_dir_path()) {
338        Ok(entries) => entries,
339        Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()),
340        Err(error) => return Err(error),
341    };
342    let mut found = Vec::new();
343    for entry in entries {
344        let path = entry?.path();
345        if !path
346            .to_str()
347            .is_some_and(|name| name.ends_with(REGISTRATION_SUFFIX))
348        {
349            continue;
350        }
351        if let Ok(Some(reg)) = read_registration_at(&path) {
352            found.push((path, reg));
353        }
354    }
355    found.sort_by(|left, right| left.0.cmp(&right.0));
356    Ok(found)
357}
358
359fn write_registration_at(path: &Path, reg: &RegistrationFile) -> io::Result<()> {
360    let json = serde_json::to_vec_pretty(reg)
361        .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
362    write_private_file(path, &json)
363}
364
365fn read_registration_at(path: &Path) -> io::Result<Option<RegistrationFile>> {
366    match std::fs::read(path) {
367        Ok(bytes) => match serde_json::from_slice(&bytes) {
368            Ok(reg) => Ok(Some(reg)),
369            Err(error) => Err(io::Error::new(
370                io::ErrorKind::InvalidData,
371                format!("{}: {error}", path.display()),
372            )),
373        },
374        Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(None),
375        Err(error) => Err(error),
376    }
377}
378
379/// Replace `path` with `bytes`, `0600`, through a temporary name in the same
380/// directory so a reader never sees half a registration.
381fn write_private_file(path: &Path, bytes: &[u8]) -> io::Result<()> {
382    use std::io::Write as _;
383    let name = path
384        .file_name()
385        .and_then(|name| name.to_str())
386        .unwrap_or("onlyne-registration");
387    let temp = path.with_file_name(format!("{name}.tmp-{}", std::process::id()));
388    let outcome = open_private(&temp).and_then(|mut file| {
389        file.write_all(bytes)?;
390        drop(file);
391        std::fs::rename(&temp, path)
392    });
393    if outcome.is_err() {
394        let _ = std::fs::remove_file(&temp);
395    }
396    outcome
397}
398
399/// Open `path` for replacement with mode `0600`.
400#[cfg(unix)]
401fn open_private(path: &Path) -> io::Result<std::fs::File> {
402    use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
403    let file = std::fs::OpenOptions::new()
404        .write(true)
405        .create(true)
406        .truncate(true)
407        .mode(0o600)
408        .open(path)?;
409    // `mode` applies at creation and passes through umask, and an existing file
410    // keeps its own mode, so the handle is chmod'd directly: a replacement is
411    // never briefly readable by group or other.
412    let mut permissions = file.metadata()?.permissions();
413    permissions.set_mode(0o600);
414    file.set_permissions(permissions)?;
415    Ok(file)
416}
417
418#[cfg(not(unix))]
419fn open_private(path: &Path) -> io::Result<std::fs::File> {
420    std::fs::OpenOptions::new()
421        .write(true)
422        .create(true)
423        .truncate(true)
424        .open(path)
425}
426
427/// The owner tree's local socket: the path v2 binds, the canonical spelling v1
428/// bound, the registration that names it, and the root all three derive from.
429///
430/// [`resolve`](Self::resolve) touches no filesystem, so a client that starts
431/// before its daemon still resolves the one path the daemon will bind, and an
432/// absent runtime directory is not an error until a caller binds or publishes.
433#[derive(Debug, Clone, PartialEq, Eq)]
434pub struct SocketEndpoint {
435    root: PathBuf,
436    natural: PathBuf,
437    actual: PathBuf,
438    registration: PathBuf,
439}
440
441impl SocketEndpoint {
442    /// Resolve from the owner root and its runtime directory
443    /// (`<root>/.onlyne/run`).
444    ///
445    /// The digest covers [`absolute_path`] of `root`, so the spelling has to
446    /// stop moving before a daemon and its clients can agree. A tree that exists
447    /// resolves its canonical spelling; one that does not yet own a directory
448    /// (a caller resolving before `bootstrap` creates it) resolves a collapsed
449    /// lexical spelling, and macOS reaches one tree through `/var` and
450    /// `/private/var`. Resolve after `bootstrap` when the derived name has to
451    /// match a running daemon's.
452    ///
453    /// `run_dir` names only [`natural`](Self::natural).
454    pub fn resolve(root: &Path, run_dir: &Path) -> Self {
455        let root = absolute_path(root);
456        let dir = runtime_dir_path();
457        let digest = workspace_digest(&root);
458        Self {
459            natural: run_dir.join(SOCKET_FILE_NAME),
460            actual: dir.join(format!("{digest}{SOCKET_SUFFIX}")),
461            registration: dir.join(format!("{digest}{REGISTRATION_SUFFIX}")),
462            root,
463        }
464    }
465
466    /// The canonical owner root the other three paths belong to.
467    pub fn root(&self) -> &Path {
468        &self.root
469    }
470
471    /// The canonical v1 spelling, `<run_dir>/s`, kept for operators and older
472    /// tooling that print it. Nothing binds or connects here under v2.
473    pub fn natural(&self) -> &Path {
474        &self.natural
475    }
476
477    /// The path `bind` and `connect` use:
478    /// `<runtime_dir>/<digest>.sock`.
479    pub fn actual(&self) -> &Path {
480        &self.actual
481    }
482
483    /// `<runtime_dir>/<digest>.json`, the file naming [`actual`](Self::actual).
484    pub fn registration(&self) -> &Path {
485        &self.registration
486    }
487
488    /// `true` when the bound path is not the canonical `<run_dir>/s` spelling,
489    /// which under v2 is every tree unless a runtime path happens to equal it.
490    /// Callers log one line about the move.
491    pub fn short(&self) -> bool {
492        self.actual != self.natural
493    }
494
495    /// Write `reg` at [`registration`](Self::registration).
496    ///
497    /// The caller owns the file: a role workspace writes its own kind, role, and
498    /// runtime here, and the next bind overwrites whatever was there.
499    pub fn publish(&self, reg: &RegistrationFile) -> io::Result<()> {
500        runtime_dir()?;
501        write_registration_at(&self.registration, reg)
502    }
503}
504
505/// Absolute, symlink-free spelling of a path when it exists, its lexical
506/// absolute otherwise.
507///
508/// `std::path::absolute` is purely lexical: it follows no symlink and keeps every
509/// `..` it is handed. macOS reaches one temporary tree through `/var` and
510/// `/private/var`, so a digest taken over a lexical spelling would split that
511/// owner tree across two derived socket directories. Canonicalizing first gives
512/// every live tree one spelling. A missing path has nothing to canonicalize, so
513/// the fallback collapses `..` against the components in front of them, which
514/// keeps `/srv/a/../b` and `/srv/b` one answer before either directory exists.
515pub fn absolute_path(path: &Path) -> PathBuf {
516    if let Ok(canonical) = std::fs::canonicalize(path) {
517        return canonical;
518    }
519    collapse_parents(&lexical_absolute(path))
520}
521
522/// Lexically cancel `..` against a preceding normal component.
523///
524/// A `..` with nothing to cancel keeps its place, so a path that climbs past its
525/// own root still spells one thing every time.
526fn collapse_parents(path: &Path) -> PathBuf {
527    let mut out = PathBuf::new();
528    for component in path.components() {
529        match component {
530            std::path::Component::ParentDir => match out.components().next_back() {
531                Some(std::path::Component::Normal(_)) => {
532                    out.pop();
533                }
534                _ => out.push(component.as_os_str()),
535            },
536            other => out.push(other.as_os_str()),
537        }
538    }
539    out
540}
541
542/// Create a directory, optionally setting its mode.
543///
544/// Shared with the layout module of `onlyne-config`, so a caller creates `run/`
545/// with the same `0700` a `bootstrap` creates, keeping one convention for the
546/// directories that hold private endpoints.
547pub fn create_dir(path: &Path, mode: Option<u32>) -> io::Result<()> {
548    std::fs::create_dir_all(path)?;
549    if let Some(mode) = mode {
550        set_dir_mode(path, mode)?;
551    }
552    Ok(())
553}
554
555#[cfg(unix)]
556pub fn set_dir_mode(path: &Path, mode: u32) -> io::Result<()> {
557    use std::os::unix::fs::PermissionsExt;
558    let mut permissions = std::fs::metadata(path)?.permissions();
559    permissions.set_mode(mode);
560    std::fs::set_permissions(path, permissions)
561}
562
563#[cfg(not(unix))]
564pub fn set_dir_mode(_path: &Path, _mode: u32) -> io::Result<()> {
565    Ok(())
566}
567
568fn set_file_mode(path: &Path, mode: u32) -> io::Result<()> {
569    let _ = std::fs::metadata(path)?;
570    set_file_mode_impl(path, mode)
571}
572
573#[cfg(unix)]
574fn set_file_mode_impl(path: &Path, mode: u32) -> io::Result<()> {
575    use std::os::unix::fs::PermissionsExt;
576    let mut permissions = std::fs::metadata(path)?.permissions();
577    permissions.set_mode(mode);
578    std::fs::set_permissions(path, permissions)
579}
580
581#[cfg(not(unix))]
582fn set_file_mode_impl(_path: &Path, _mode: u32) -> io::Result<()> {
583    Ok(())
584}
585
586/// Tokio listener produced by [`bind_local`].
587pub type LocalListener = interprocess::local_socket::tokio::Listener;
588/// Tokio stream produced by [`connect_local`] or [`LocalListener::accept`](interprocess::local_socket::tokio::prelude::Listener::accept).
589pub type LocalStream = interprocess::local_socket::tokio::Stream;
590/// Blocking listener for tests that cannot run inside a tokio runtime.
591pub type LocalListenerSync = interprocess::local_socket::Listener;
592/// Blocking stream matching [`LocalListenerSync`].
593pub type LocalStreamSync = interprocess::local_socket::Stream;
594
595/// Traits needed to `.accept()` / `.connect()` the interprocess types.
596pub mod prelude {
597    pub use interprocess::local_socket::traits::tokio::{
598        Listener as TokioListener, Stream as TokioStream,
599    };
600    pub use interprocess::local_socket::traits::{Listener as SyncListener, Stream as SyncStream};
601}
602
603#[cfg(windows)]
604const MARKER_PREFIX: &str = "v1:";
605const PIPE_BUSY: i32 = 231;
606const VERBATIM_PIPE_PREFIX: &str = r"\\.\pipe\";
607
608/// `true` when `--socket` is already an NPFS path and must not be hashed.
609pub fn is_verbatim_pipe_path(path: &Path) -> bool {
610    path.to_str()
611        .is_some_and(|s| starts_with_ignore_ascii_case(s, VERBATIM_PIPE_PREFIX))
612}
613
614/// Lexical absolute spelling of `path`.
615///
616/// `std::path::absolute` (Rust 1.79+) never touches the filesystem, so a missing
617/// path still gets an answer.
618fn lexical_absolute(path: &Path) -> PathBuf {
619    std::path::absolute(path).unwrap_or_else(|_| path.to_path_buf())
620}
621
622/// Lowercase hex of the first `bytes` of `sha256` over `path`.
623///
624/// Separators become `/` and the whole string is lowercased, so `C:\Work` and
625/// `c:/work` digest alike.
626fn digest_hex(path: &Path, bytes: usize) -> String {
627    let normalized = path
628        .to_string_lossy()
629        .replace('\\', "/")
630        .to_ascii_lowercase();
631    let digest = Sha256::digest(normalized.as_bytes());
632    let mut hex = String::with_capacity(bytes * 2);
633    for byte in &digest[..bytes] {
634        hex.push(HEX[(*byte >> 4) as usize] as char);
635        hex.push(HEX[(*byte & 0x0f) as usize] as char);
636    }
637    hex
638}
639
640/// NPFS leaf `onlyne-<32hex>` derived from `path` without touching the filesystem.
641///
642/// `std::path::absolute` is lexical. Separators become `/` and the whole string
643/// is lowercased before the digest so `C:\Work` and `c:/work` agree.
644pub fn pipe_name_for(path: &Path) -> String {
645    format!("onlyne-{}", digest_hex(&lexical_absolute(path), 16))
646}
647
648/// Identity of one owner tree: `sha256` of its canonical spelling, the first 16
649/// lowercase hex characters.
650///
651/// The spelling is [`absolute_path`], so macOS's `/var` and `/private/var`,
652/// `a/../b` and `b`, and `C:\Work` and `c:/work` each digest alike — the whole
653/// reason a daemon and every client can derive one file name without talking to
654/// each other. Both runtime files of a tree share this digest, one leaf apart.
655pub fn workspace_digest(root: &Path) -> String {
656    digest_hex(&absolute_path(root), 8)
657}
658
659const HEX: &[u8; 16] = b"0123456789abcdef";
660
661/// Bind a tokio listener at `path`.
662///
663/// Callers still remove a stale `path` first. Unix `mode(0o600)` is set on the
664/// bind options (fchmod before bind, no umask TOCTOU). Windows writes the
665/// marker, then binds the derived pipe with an owner-only SDDL; a live listener
666/// holding that NPFS name makes the second bind fail, which is the EADDRINUSE
667/// equivalent (`try_overwrite` is a no-op on Windows).
668#[allow(clippy::unused_async)]
669pub async fn bind_local(path: &Path) -> io::Result<LocalListener> {
670    bind_tokio(path)
671}
672
673/// Synchronous counterpart of [`bind_local`] for the tokio listener type.
674pub fn bind_tokio(path: &Path) -> io::Result<LocalListener> {
675    create_with_privacy(
676        path,
677        ListenerNonblockingMode::Neither,
678        ListenerOptions::create_tokio,
679    )
680}
681
682/// Bind the owner tree's socket in the runtime directory.
683///
684/// The runtime directory is created and owner-verified, a stale name at the
685/// bound path is dropped, and nothing else is written: the registration naming
686/// this socket is the serving side's to publish, because only the serving side
687/// knows whether the tree is an admin root or a role workspace
688/// ([`write_registration`], [`SocketEndpoint::publish`]).
689pub fn bind_socket_v2(root: &Path) -> io::Result<LocalListener> {
690    let path = socket_path(root)?;
691    bind_runtime_socket(&path)
692}
693
694/// Bind the owner tree's local socket at [`socket_path`] in the runtime
695/// directory and return it with the endpoint that names it.
696///
697/// `root` is the owner tree's root and `run_dir` is `<root>/.onlyne/run`, which
698/// is created `0700` `bootstrap`-style for a caller that binds before
699/// `bootstrap` has, and which survives only as [`SocketEndpoint::natural`].
700/// Resolution happens before that creation, so a caller that needs binding and a
701/// client to agree on the digest resolves after `bootstrap`, which guarantees a
702/// canonical root.
703pub fn bind_socket(root: &Path, run_dir: &Path) -> io::Result<(LocalListener, SocketEndpoint)> {
704    let endpoint = SocketEndpoint::resolve(root, run_dir);
705    runtime_dir()?;
706    create_dir(run_dir, Some(0o700))?;
707    let listener = bind_runtime_socket(endpoint.actual())?;
708    Ok((listener, endpoint))
709}
710
711/// Bind the owner tree's socket and publish `reg` for it in one call, for a
712/// daemon that knows its surface at start.
713///
714/// A refused registration fails the call: a socket no reader can find is not the
715/// outcome the serving side asked for.
716pub fn bind_socket_registered(
717    root: &Path,
718    run_dir: &Path,
719    reg: &RegistrationFile,
720) -> io::Result<(LocalListener, SocketEndpoint)> {
721    let (listener, endpoint) = bind_socket(root, run_dir)?;
722    endpoint.publish(reg)?;
723    Ok((listener, endpoint))
724}
725
726/// Drop a stale name at `path` and bind it, naming the path and its runtime
727/// directory in any failure.
728fn bind_runtime_socket(path: &Path) -> io::Result<LocalListener> {
729    #[cfg(unix)]
730    remove_stale_socket(path);
731    bind_tokio(path).map_err(|error| bind_failure(path, error))
732}
733
734/// Unix: drop the name so a restarted daemon can take it again.
735///
736/// `reclaim_name(false)` leaves unlink to the owner, and a name left in place
737/// makes the bind below fail with the path and its length in the message.
738#[cfg(unix)]
739fn remove_stale_socket(path: &Path) {
740    let _ = std::fs::remove_file(path);
741}
742
743/// A bind failure the operator can act on: the path that was tried, its length,
744/// and the runtime directory holding it.
745fn bind_failure(path: &Path, error: io::Error) -> io::Error {
746    let dir = path.parent().unwrap_or(Path::new(""));
747    io::Error::new(
748        error.kind(),
749        format!(
750            "bind {} ({} bytes) in {}: {error}",
751            path.display(),
752            path.as_os_str().len(),
753            dir.display(),
754        ),
755    )
756}
757
758/// Adopt a directory this process owns outright.
759///
760/// Ownership is the check that decides. `/tmp` is world-executable, so a
761/// pre-created `/tmp/onlyne-<uid>` owned by another user is refused by uid:
762/// a socket inside a directory that user can rewrite is one that user could
763/// replace under a live daemon.
764///
765/// A directory this user owns that is open to group or other is tightened to
766/// `0700` rather than refused, because the endpoints inside it are private only
767/// while the directory is, and its mode is whatever the first writer chose —
768/// an installer, a `tmpfiles.d` rule, or a client that made the directory
769/// before its daemon did. Refusing would brick every later daemon over a mode
770/// that carries no ownership information, which is why the uid check above and
771/// not the mode is the decision.
772#[cfg(unix)]
773fn ensure_private_dir(path: &Path) -> io::Result<()> {
774    match std::fs::symlink_metadata(path) {
775        Ok(metadata) => verify_private_dir(path, &metadata)?,
776        Err(error) if error.kind() == io::ErrorKind::NotFound => {
777            // `create_dir_all` because an operator-named runtime directory can
778            // nest under paths that do not exist yet, which the default
779            // `/tmp/onlyne-<uid>` never does. A process that loses the race to
780            // create it gets `Ok` and the chmod below, which fails with EPERM
781            // when the winner was another user and no-ops when it was this one.
782            std::fs::create_dir_all(path).map_err(|source| dir_failure(path, &source))?;
783            set_dir_mode(path, 0o700).map_err(|source| dir_failure(path, &source))?;
784        }
785        Err(source) => return Err(dir_failure(path, &source)),
786    }
787    // Mode `0700` is open to exactly one uid, and the file system owner can
788    // write anywhere (and macOS ACLs can deny that owner), so the create-new
789    // probe is the decision this process actually makes. The name carries a
790    // sequence number as well as the pid, so two threads of one process
791    // adopting the same directory do not collide on the probe itself.
792    static PROBE_SEQUENCE: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
793    let sequence = PROBE_SEQUENCE.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
794    let probe = path.join(format!("onlyne-probe-{}-{sequence}", std::process::id()));
795    match std::fs::OpenOptions::new()
796        .write(true)
797        .create_new(true)
798        .open(&probe)
799    {
800        Ok(_) => {
801            let _ = std::fs::remove_file(&probe);
802            Ok(())
803        }
804        Err(source) => Err(dir_failure(path, &source)),
805    }
806}
807
808/// Refuse anything but a directory this user owns and no one else can enter.
809#[cfg(unix)]
810fn verify_private_dir(path: &Path, metadata: &std::fs::Metadata) -> io::Result<()> {
811    use std::os::unix::fs::{MetadataExt, PermissionsExt};
812    if metadata.file_type().is_symlink() {
813        return Err(dir_refusal(path, "a symlink"));
814    }
815    if !metadata.is_dir() {
816        return Err(dir_refusal(path, "already held by a non-directory"));
817    }
818    if metadata.uid() != unsafe { libc::geteuid() } {
819        return Err(dir_refusal(path, "owned by another user"));
820    }
821    if metadata.permissions().mode() & 0o077 != 0 {
822        set_dir_mode(path, 0o700).map_err(|source| dir_failure(path, &source))?;
823    }
824    Ok(())
825}
826
827/// Off unix the temporary root's ACLs are the guard, so the runtime directory is
828/// created if absent and otherwise left alone.
829#[cfg(not(unix))]
830fn ensure_private_dir(path: &Path) -> io::Result<()> {
831    if path.is_dir() {
832        return Ok(());
833    }
834    std::fs::create_dir_all(path)
835}
836
837#[cfg(unix)]
838fn dir_refusal(path: &Path, reason: &str) -> io::Error {
839    io::Error::new(
840        io::ErrorKind::PermissionDenied,
841        format!("refusing to use directory {}: {reason}", path.display()),
842    )
843}
844
845#[cfg(unix)]
846fn dir_failure(path: &Path, source: &io::Error) -> io::Error {
847    io::Error::new(
848        source.kind(),
849        format!("directory {}: {source}", path.display()),
850    )
851}
852
853/// Blocking bind, used by CLI tests that serve one frame on a helper thread.
854pub fn bind_local_sync(path: &Path) -> io::Result<LocalListenerSync> {
855    create_with_privacy(
856        path,
857        ListenerNonblockingMode::Neither,
858        ListenerOptions::create_sync,
859    )
860}
861
862/// Blocking bind whose `accept` returns `WouldBlock` when no client is waiting.
863pub fn bind_local_sync_poll(path: &Path) -> io::Result<LocalListenerSync> {
864    create_with_privacy(
865        path,
866        ListenerNonblockingMode::Accept,
867        ListenerOptions::create_sync,
868    )
869}
870
871/// Connect to the listener that [`bind_local`] created for `path`.
872///
873/// Windows `ERROR_PIPE_BUSY` (231) is remapped to [`io::ErrorKind::WouldBlock`]
874/// so a caller can retry inside its existing timeout budget.
875pub async fn connect_local(path: &Path) -> io::Result<LocalStream> {
876    use interprocess::local_socket::tokio::Stream;
877    use interprocess::local_socket::traits::tokio::Stream as _;
878    Stream::connect(connect_name(path)?)
879        .await
880        .map_err(map_pipe_busy)
881}
882
883/// Blocking connect matching [`bind_local_sync`].
884pub fn connect_local_sync(path: &Path) -> io::Result<LocalStreamSync> {
885    use interprocess::local_socket::Stream;
886    use interprocess::local_socket::traits::Stream as _;
887    Stream::connect(connect_name(path)?).map_err(map_pipe_busy)
888}
889
890fn map_pipe_busy(err: io::Error) -> io::Error {
891    if err.raw_os_error() == Some(PIPE_BUSY) {
892        io::Error::new(io::ErrorKind::WouldBlock, err)
893    } else {
894        err
895    }
896}
897
898fn starts_with_ignore_ascii_case(value: &str, prefix: &str) -> bool {
899    value.len() >= prefix.len()
900        && value
901            .as_bytes()
902            .iter()
903            .zip(prefix.as_bytes())
904            .all(|(a, b)| a.eq_ignore_ascii_case(b))
905}
906
907#[cfg(windows)]
908fn parse_marker(text: &str) -> Option<String> {
909    let rest = text.trim().strip_prefix(MARKER_PREFIX)?;
910    if rest.is_empty()
911        || !rest
912            .bytes()
913            .all(|b| b.is_ascii_graphic() && b != b'/' && b != b'\\')
914    {
915        return None;
916    }
917    Some(rest.to_string())
918}
919
920#[cfg(windows)]
921fn read_marker_name(path: &Path) -> Option<String> {
922    let text = std::fs::read_to_string(path).ok()?;
923    parse_marker(&text)
924}
925
926#[cfg(windows)]
927fn write_marker(path: &Path, pipe_name: &str) -> io::Result<()> {
928    std::fs::write(path, format!("{MARKER_PREFIX}{pipe_name}"))
929}
930
931#[cfg(windows)]
932fn windows_pipe_leaf(path: &Path) -> String {
933    read_marker_name(path).unwrap_or_else(|| pipe_name_for(path))
934}
935
936fn create_with_privacy<T>(
937    path: &Path,
938    nonblocking: ListenerNonblockingMode,
939    create: fn(ListenerOptions<'static>) -> io::Result<T>,
940) -> io::Result<T> {
941    #[cfg(unix)]
942    {
943        use interprocess::os::unix::local_socket::ListenerOptionsExt;
944        match create(unix_options(path, nonblocking)?.mode(0o600)) {
945            Ok(listener) => Ok(listener),
946            Err(err) if err.kind() == io::ErrorKind::Unsupported => {
947                // macOS: fchmod on an unbound socket is Unsupported. chmod
948                // after bind is the 1.0.x path and still yields 0600.
949                let listener = create(unix_options(path, nonblocking)?)?;
950                apply_private_mode(path).map_err(|e| io::Error::other(e.to_string()))?;
951                Ok(listener)
952            }
953            Err(err) => Err(err),
954        }
955    }
956    #[cfg(windows)]
957    {
958        create(listener_options_windows(path, nonblocking)?)
959    }
960}
961
962#[cfg(unix)]
963fn unix_options(
964    path: &Path,
965    nonblocking: ListenerNonblockingMode,
966) -> io::Result<ListenerOptions<'static>> {
967    let name = path.to_fs_name::<GenericFilePath>()?.into_owned();
968    Ok(ListenerOptions::new()
969        .name(name)
970        .reclaim_name(false)
971        .nonblocking(nonblocking))
972}
973
974#[cfg(windows)]
975fn listener_options_windows(
976    path: &Path,
977    nonblocking: ListenerNonblockingMode,
978) -> io::Result<ListenerOptions<'static>> {
979    use interprocess::local_socket::{GenericNamespaced, ToNsName};
980    use interprocess::os::windows::local_socket::ListenerOptionsExt;
981    use interprocess::os::windows::security_descriptor::SecurityDescriptor;
982    use widestring::u16cstr;
983
984    let name = if is_verbatim_pipe_path(path) {
985        path.to_fs_name::<GenericFilePath>()?.into_owned()
986    } else {
987        let leaf = windows_pipe_leaf(path);
988        write_marker(path, &leaf)?;
989        leaf.to_ns_name::<GenericNamespaced>()?.into_owned()
990    };
991    let sd = SecurityDescriptor::deserialize(u16cstr!("D:P(A;;GA;;;OW)(A;;GA;;;SY)"))?;
992    Ok(ListenerOptions::new()
993        .name(name)
994        .reclaim_name(false)
995        .nonblocking(nonblocking)
996        .security_descriptor(sd))
997}
998
999fn connect_name(path: &Path) -> io::Result<interprocess::local_socket::Name<'static>> {
1000    #[cfg(unix)]
1001    {
1002        Ok(path.to_fs_name::<GenericFilePath>()?.into_owned())
1003    }
1004    #[cfg(windows)]
1005    {
1006        use interprocess::local_socket::{GenericNamespaced, ToNsName};
1007        if is_verbatim_pipe_path(path) {
1008            Ok(path.to_fs_name::<GenericFilePath>()?.into_owned())
1009        } else {
1010            Ok(windows_pipe_leaf(path)
1011                .to_ns_name::<GenericNamespaced>()?
1012                .into_owned())
1013        }
1014    }
1015}