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(®istration_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(®.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}