diff --git a/kernel/DAEMON-LIFECYCLE.md b/kernel/DAEMON-LIFECYCLE.md index fc5dbaf28..3d739fab2 100644 --- a/kernel/DAEMON-LIFECYCLE.md +++ b/kernel/DAEMON-LIFECYCLE.md @@ -31,8 +31,17 @@ ambiguous case below. A file is never believed on its own. ### Path -`/connection.json` — `data-dir` is the same directory that already -holds `relayflowd.sock` and `relayflowd.sqlite3`, default `.relayflowd`. +`/connection.json`, default data-dir `.relayflowd`. The data dir also +holds `relayflowd.sqlite3`, `relayflowd.lock`, and `relayflowd.log`. + +**The socket itself lives OUTSIDE the data dir**, at a short hashed path +under `$XDG_RUNTIME_DIR` / `$TMPDIR` — see `kernel/relayflowd/src/socket_path.rs` +and `socketPathFor` in `packages/sdk/src/daemon-connection.ts`. The daemon and +the CLI derive the same path from the same absolute data-dir input, so §2 +step 3's lexical equality still holds. This decoupling exists so a deep +working directory cannot push the full socket path past `SUN_LEN` (~104 bytes +on macOS), which used to fail the daemon at `bind(2)` before any step could +run (#262). ### Shape diff --git a/kernel/relayflowd/src/lib.rs b/kernel/relayflowd/src/lib.rs index 6c964cb8d..85e3fe674 100644 --- a/kernel/relayflowd/src/lib.rs +++ b/kernel/relayflowd/src/lib.rs @@ -3,6 +3,7 @@ pub mod engine; pub mod exec_det; pub mod memory; pub mod server; +pub mod socket_path; pub mod worker; pub use engine::{ diff --git a/kernel/relayflowd/src/server.rs b/kernel/relayflowd/src/server.rs index 95a42df99..68037205b 100644 --- a/kernel/relayflowd/src/server.rs +++ b/kernel/relayflowd/src/server.rs @@ -59,11 +59,14 @@ pub fn serve(data_dir: &Path) -> Result<()> { } Err(lifecycle::AcquireError::Io(error)) => return Err(error.into()), }; - let socket_path = data_dir.join("relayflowd.sock"); - lifecycle::remove_residue(data_dir)?; + // Bind outside the data dir at a short hashed path so SUN_LEN cannot be + // violated by a deep working directory (#262). See socket_path.rs and + // DAEMON-LIFECYCLE.md §1. + let socket_path = crate::socket_path::derive_socket_path(data_dir)?; + lifecycle::remove_residue(data_dir, &socket_path)?; let listener = UnixListener::bind(&socket_path) .with_context(|| format!("bind socket {}", socket_path.display()))?; - lifecycle::publish(&socket_path)?; + lifecycle::publish(data_dir, &socket_path)?; let hub = Arc::new(ProtocolHub::default()); reconcile::spawn_reconciler(data_dir.to_path_buf(), hub.clone()); // Trigger-plane liveness sweep (RFC-0001 gate 2, Native silent-death diff --git a/kernel/relayflowd/src/server/client.rs b/kernel/relayflowd/src/server/client.rs index b8ff5a846..3e3fe7b03 100644 --- a/kernel/relayflowd/src/server/client.rs +++ b/kernel/relayflowd/src/server/client.rs @@ -22,7 +22,7 @@ fn lifecycle_request_via_socket( io::{BufRead, BufReader, Write}, os::unix::net::UnixStream, }; - let socket = data_dir.join("relayflowd.sock"); + let socket = crate::socket_path::derive_socket_path(data_dir)?; if !socket.exists() { return Ok(None); } @@ -110,7 +110,7 @@ pub fn resume_via_socket(data_dir: &Path, run_id: &str) -> Result std::result::Result Result<()> { - for name in ["relayflowd.sock", "connection.json"] { - let path = data_dir.join(name); - match std::fs::remove_file(&path) { +/// +/// `socket_path` is now outside the data dir (see socket_path.rs and +/// DAEMON-LIFECYCLE.md §1), so it is passed explicitly rather than derived +/// from the data dir alone. The lock proves nobody else owns it. +pub(crate) fn remove_residue(data_dir: &Path, socket_path: &Path) -> Result<()> { + let connection = data_dir.join("connection.json"); + for path in [socket_path, connection.as_path()] { + match std::fs::remove_file(path) { Ok(()) => {} Err(error) if error.kind() == io::ErrorKind::NotFound => {} Err(error) => { @@ -86,7 +90,12 @@ struct ConnectionFile<'a> { } /// Publish only after `UnixListener::bind`, which has already called listen(2). -pub(crate) fn publish(socket_path: &Path) -> Result<()> { +/// +/// `data_dir` and `socket_path` used to be one and the same directory +/// (`connection.json` sat next to `relayflowd.sock`). #262 moves the socket +/// outside the data dir so they diverge; the connection file still lives in +/// the data dir, alongside the lock and journal. +pub(crate) fn publish(data_dir: &Path, socket_path: &Path) -> Result<()> { // `absolute` preserves the caller's data-dir spelling (unlike // `canonicalize`, which would turn `/var` into macOS's `/private/var` and // fail the client's lexical `resolve(dataDir)` comparison). @@ -95,7 +104,6 @@ pub(crate) fn publish(socket_path: &Path) -> Result<()> { let socket_text = socket_path .to_str() .context("socket path is not valid UTF-8")?; - let data_dir = socket_path.parent().context("socket path has no parent")?; let connection_path = data_dir.join("connection.json"); let temporary_path = data_dir.join(format!("connection.json.tmp.{}", std::process::id())); let record = ConnectionFile { diff --git a/kernel/relayflowd/src/socket_path.rs b/kernel/relayflowd/src/socket_path.rs new file mode 100644 index 000000000..36667110d --- /dev/null +++ b/kernel/relayflowd/src/socket_path.rs @@ -0,0 +1,126 @@ +//! Where the unix-domain socket lives, and why not inside the data dir. +//! +//! `/relayflowd.sock` was the naive location, and it fails hard when +//! the working directory pushes the full socket path past the OS `SUN_LEN` +//! limit (~104 bytes on macOS, ~108 on Linux). A user cloning an example into +//! a nested tree hits `Error: bind socket ...: path must be shorter than +//! SUN_LEN` before any step can run (#262). +//! +//! The daemon socket now binds under `$XDG_RUNTIME_DIR` / `$TMPDIR` / +//! `std::env::temp_dir()` at a short, stable path keyed by a 12-hex-char +//! SHA-256 of the *absolute* data-dir path. Anti-hijack still holds because +//! both the daemon and the CLI derive the same path from the same input: +//! [`derive_socket_path`] is the Rust side, `socketPathFor` in +//! `packages/sdk/src/daemon-connection.ts` is the identical JS side. The +//! `socket_path` field in `connection.json` still names the bound socket, and +//! the client's §2 step-3 check still refuses a file that names anything else. +//! +//! The lock (`relayflowd.lock`), the journal (`relayflowd.sqlite3`), the log +//! (`relayflowd.log`), and `connection.json` continue to live in the data dir +//! — nothing there has a length problem, and moving them would break the +//! `flock(2)`-based mutex that DAEMON-LIFECYCLE.md §3 depends on. + +use std::path::{Path, PathBuf}; + +use anyhow::{Context, Result}; +use sha2::{Digest, Sha256}; + +/// Derive the socket path from the data dir. +/// +/// Uses `std::path::absolute` — never `canonicalize` — so the JS side's +/// lexical `path.resolve(dataDir)` produces the same input string. Symlink +/// resolution would introduce a divergence (macOS's `/var` → `/private/var` +/// is the classic one), and the client validates by string equality. +pub fn derive_socket_path(data_dir: &Path) -> Result { + let absolute = std::path::absolute(data_dir) + .with_context(|| format!("make data dir absolute {}", data_dir.display()))?; + let hash = hash_data_dir(&absolute); + Ok(runtime_dir().join(format!("relayflowd-{hash}.sock"))) +} + +fn hash_data_dir(absolute_data_dir: &Path) -> String { + let mut hasher = Sha256::new(); + hasher.update(absolute_data_dir.as_os_str().as_encoded_bytes()); + let digest = hasher.finalize(); + let mut out = String::with_capacity(12); + for byte in &digest[..6] { + use std::fmt::Write; + // Fixed-width lowercase hex, matching Node's `digest('hex').slice(0, 12)`. + write!(&mut out, "{byte:02x}").expect("write to String is infallible"); + } + out +} + +fn runtime_dir() -> PathBuf { + // XDG first (Linux, typically `/run/user//`, ~14 chars — the shortest + // per-user private option). TMPDIR next (macOS, `/var/folders/xx/YYY/T/`, + // ~50 chars, still per-user private). Only fall back to the system-wide + // `std::env::temp_dir()` when neither is set; on macOS that resolves to + // `/var/folders/.../T` too, and on Linux to `/tmp` which is world-writable. + if let Some(value) = env_nonempty("XDG_RUNTIME_DIR") { + return PathBuf::from(value); + } + if let Some(value) = env_nonempty("TMPDIR") { + return PathBuf::from(value); + } + std::env::temp_dir() +} + +fn env_nonempty(name: &str) -> Option { + match std::env::var_os(name) { + Some(value) if !value.is_empty() => Some(value), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn same_data_dir_yields_same_socket() { + let a = derive_socket_path(Path::new("/tmp/deep")).unwrap(); + let b = derive_socket_path(Path::new("/tmp/deep")).unwrap(); + assert_eq!(a, b); + } + + #[test] + fn different_data_dirs_yield_different_sockets() { + let a = derive_socket_path(Path::new("/tmp/one")).unwrap(); + let b = derive_socket_path(Path::new("/tmp/two")).unwrap(); + assert_ne!(a, b); + } + + #[test] + fn deep_data_dir_produces_short_socket_path() { + // The whole point of #262: the socket path stays short regardless of + // how deep the data dir is. Assert against a padding well past the + // 104-byte macOS SUN_LEN. + let deep = "/tmp/".to_string() + &"a/".repeat(80) + "data"; + assert!(deep.len() > 104); + let socket = derive_socket_path(Path::new(&deep)).unwrap(); + // The socket path is bounded by ` + "/relayflowd-" + + // 12 hex chars + ".sock"`, which is ~60 chars on macOS worst case. + assert!( + socket.as_os_str().len() < 104, + "derived socket path must fit SUN_LEN, got {} bytes: {}", + socket.as_os_str().len(), + socket.display(), + ); + } + + #[test] + fn relative_and_absolute_data_dirs_agree() { + // A caller passing a relative `--data-dir` must produce the same + // socket as a caller passing the same absolute path, or the CLI's + // lexical anti-hijack check would refuse a matching daemon. + // `absolute` resolves relative paths against the process CWD, and + // that is exactly what Node's `path.resolve` does on the other side. + let cwd = std::env::current_dir().unwrap(); + let absolute = cwd.join("some-data"); + assert_eq!( + derive_socket_path(Path::new("some-data")).unwrap(), + derive_socket_path(&absolute).unwrap(), + ); + } +} diff --git a/kernel/relayflowd/tests/crash_resume.rs b/kernel/relayflowd/tests/crash_resume.rs index c77613d83..8b51e80a4 100644 --- a/kernel/relayflowd/tests/crash_resume.rs +++ b/kernel/relayflowd/tests/crash_resume.rs @@ -124,7 +124,7 @@ fn sigkill_under_serve_resumes_the_socket_started_run() { .process_group(0) .spawn() .unwrap(); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); wait_until("serve socket", || socket.exists()); let spec: Value = serde_json::from_slice(&fs::read(&fixture.spec_path).unwrap()).unwrap(); diff --git a/kernel/relayflowd/tests/crash_resume/agent_support.rs b/kernel/relayflowd/tests/crash_resume/agent_support.rs index b17b68d16..ae9386a3a 100644 --- a/kernel/relayflowd/tests/crash_resume/agent_support.rs +++ b/kernel/relayflowd/tests/crash_resume/agent_support.rs @@ -100,7 +100,8 @@ impl AgentFixture { } pub fn socket(&self) -> PathBuf { - self.data_dir.join("relayflowd.sock") + relayflowd::socket_path::derive_socket_path(&self.data_dir) + .expect("derive socket path") } pub fn server(&self) -> ServerGuard { diff --git a/kernel/relayflowd/tests/crash_resume/concurrency.rs b/kernel/relayflowd/tests/crash_resume/concurrency.rs index b3e51c7e4..81a5f9c84 100644 --- a/kernel/relayflowd/tests/crash_resume/concurrency.rs +++ b/kernel/relayflowd/tests/crash_resume/concurrency.rs @@ -60,7 +60,7 @@ fn run_start_dispatches_every_independent_lane_before_any_completion() { // A live resume sees both leases held and must neither abandon nor // redispatch either attempt. - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); let resumed = control .request("run.resume", json!({"run_id": run_id})) .unwrap(); @@ -161,7 +161,7 @@ fn concurrent_resumes_lease_exactly_one_attempt() { let run_id = start_run(&fixture); let mut worker = attached_worker(&fixture, "race-stub"); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); let barrier = Arc::new(Barrier::new(2)); let resumes = (0..2) .map(|_| { @@ -231,7 +231,7 @@ fn live_resume_leaves_an_active_lease_running() { .unwrap(); assert!(heartbeat["lease_deadline_ms"].as_i64().unwrap() > 0); - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); let resumed = control .request("run.resume", json!({"run_id": run_id})) .unwrap(); @@ -266,7 +266,7 @@ fn cancel_closes_the_lease_and_rejects_a_late_completion() { let mut worker = attached_worker(&fixture, "cancel-stub"); let run_id = start_run(&fixture); let dispatch = worker.event("step.dispatch").unwrap(); - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); let canceled = control .request("run.cancel", json!({"run_id": run_id})) @@ -298,7 +298,7 @@ fn cancel_and_completion_race_has_one_terminal_fact() { let mut worker = attached_worker(&fixture, "race-stub"); let run_id = start_run(&fixture); let dispatch = worker.event("step.dispatch").unwrap(); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); let barrier = Arc::new(Barrier::new(2)); let cancel_barrier = barrier.clone(); diff --git a/kernel/relayflowd/tests/crash_resume/llm.rs b/kernel/relayflowd/tests/crash_resume/llm.rs index d0dcf2f59..ff367ef1b 100644 --- a/kernel/relayflowd/tests/crash_resume/llm.rs +++ b/kernel/relayflowd/tests/crash_resume/llm.rs @@ -22,7 +22,7 @@ fn serve_plumbs_watch_events_and_replayable_stream_verbs() { let fixture = LlmFixture::new("protocol-verbs", false); let _server = ServerGuard::start(&fixture); let run_id = start_run(&fixture); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); let mut watcher = ProtocolClient::connect(&socket); watcher .request("run.watch", json!({"run_id": run_id})) diff --git a/kernel/relayflowd/tests/crash_resume/llm_support.rs b/kernel/relayflowd/tests/crash_resume/llm_support.rs index b35975de7..01c4dbf67 100644 --- a/kernel/relayflowd/tests/crash_resume/llm_support.rs +++ b/kernel/relayflowd/tests/crash_resume/llm_support.rs @@ -190,8 +190,9 @@ impl LlmFixture { fixture } - fn socket(&self) -> PathBuf { - self.data_dir.join("relayflowd.sock") + pub fn socket(&self) -> PathBuf { + relayflowd::socket_path::derive_socket_path(&self.data_dir) + .expect("derive socket path") } } diff --git a/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs b/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs index 62f24e71f..8071130e5 100644 --- a/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs +++ b/kernel/relayflowd/tests/crash_resume/parallel_lifecycle.rs @@ -16,7 +16,7 @@ use super::{ }; fn attached_agent(fixture: &LlmFixture, id: &str) -> ProtocolClient { - let mut worker = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut worker = ProtocolClient::connect(&fixture.socket()); worker .request( "worker.attach", @@ -101,7 +101,7 @@ fn renewed_parallel_leases_survive_the_original_grant_and_remain_distinct() { thread::sleep(Duration::from_millis(40)); let lane_a_deadline = renew(&mut worker, &dispatches[1]); assert_ne!(lane_b_deadline, lane_a_deadline); - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); let snapshot = control .request("run.get", json!({"run_id": run_id})) .unwrap(); @@ -302,7 +302,7 @@ fn terminal_failure_drains_or_explains_every_live_sibling() { == 2 }) }); - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); assert_eq!( control .request("run.resume", json!({"run_id": run_id})) diff --git a/kernel/relayflowd/tests/crash_resume/pin_projection.rs b/kernel/relayflowd/tests/crash_resume/pin_projection.rs index 2e1c1c013..157cc9807 100644 --- a/kernel/relayflowd/tests/crash_resume/pin_projection.rs +++ b/kernel/relayflowd/tests/crash_resume/pin_projection.rs @@ -12,7 +12,7 @@ use super::{ fn rejected_completion_cannot_forge_inspect_retry_pins_over_the_real_socket() { let fixture = LlmFixture::parallel("rejected-pin-projection"); let _server = ServerGuard::start(&fixture); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); let mut worker = ProtocolClient::connect(&socket); worker .request( diff --git a/kernel/relayflowd/tests/crash_resume/protocol_admission.rs b/kernel/relayflowd/tests/crash_resume/protocol_admission.rs index a136fe614..8d2c15dce 100644 --- a/kernel/relayflowd/tests/crash_resume/protocol_admission.rs +++ b/kernel/relayflowd/tests/crash_resume/protocol_admission.rs @@ -12,7 +12,7 @@ use super::{ fn every_mutating_run_verb_refuses_terminal_before_changing_state() { let fixture = LlmFixture::completed("terminal-admission"); let _server = ServerGuard::start(&fixture); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); let mut client = ProtocolClient::connect(&socket); let started = client .request( diff --git a/kernel/relayflowd/tests/crash_resume/support.rs b/kernel/relayflowd/tests/crash_resume/support.rs index 80ce8e654..a6bbe0c0d 100644 --- a/kernel/relayflowd/tests/crash_resume/support.rs +++ b/kernel/relayflowd/tests/crash_resume/support.rs @@ -37,6 +37,11 @@ impl Fixture { Self::new(name, true) } + pub fn socket(&self) -> PathBuf { + relayflowd::socket_path::derive_socket_path(&self.data_dir) + .expect("derive socket path") + } + fn new(name: &str, blocking_second: bool) -> Self { let directory = tempdir().unwrap(); let data_dir = directory.path().join("data"); diff --git a/kernel/relayflowd/tests/crash_resume/surface_identity.rs b/kernel/relayflowd/tests/crash_resume/surface_identity.rs index 708dcaecc..4fb368c93 100644 --- a/kernel/relayflowd/tests/crash_resume/surface_identity.rs +++ b/kernel/relayflowd/tests/crash_resume/surface_identity.rs @@ -30,7 +30,7 @@ fn complete(worker: &mut ProtocolClient, dispatch: &Value) -> Value { fn aliases_are_rejected_and_external_ancestors_serialize_over_real_sockets() { let fixture = LlmFixture::parallel("surface-identity"); let _server = ServerGuard::start(&fixture); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); let mut worker = ProtocolClient::connect(&socket); worker .request( diff --git a/kernel/relayflowd/tests/crash_resume/worker_capacity.rs b/kernel/relayflowd/tests/crash_resume/worker_capacity.rs index 8eafc8580..682bc7af0 100644 --- a/kernel/relayflowd/tests/crash_resume/worker_capacity.rs +++ b/kernel/relayflowd/tests/crash_resume/worker_capacity.rs @@ -9,7 +9,7 @@ use super::{ }; fn attach(fixture: &LlmFixture, worker_id: &str, capacity: Option) -> ProtocolClient { - let mut worker = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut worker = ProtocolClient::connect(&fixture.socket()); let mut params = json!({"worker_id": worker_id, "step_types": ["llm"]}); if let Some(capacity) = capacity { params["capacity"] = json!(capacity); @@ -24,7 +24,7 @@ fn two_workers_receive_a_deterministic_fair_capacity_bounded_batch() { let _server = ServerGuard::start(&fixture); let mut first = attach(&fixture, "first", Some(2)); let mut second = attach(&fixture, "second", Some(2)); - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); let started = control .request( "run.start", @@ -78,7 +78,7 @@ fn default_capacity_one_reopens_only_after_durable_completion_or_crash() { let fixture = LlmFixture::parallel("capacity-one-completion"); let _server = ServerGuard::start(&fixture); let mut worker = attach(&fixture, "serial", None); - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); let started = control .request( "run.start", @@ -118,7 +118,7 @@ fn default_capacity_one_reopens_only_after_durable_completion_or_crash() { let fixture = LlmFixture::parallel("capacity-one-crash"); let _server = ServerGuard::start(&fixture); let mut crashed = attach(&fixture, "crashed", None); - let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + let mut control = ProtocolClient::connect(&fixture.socket()); let started = control .request( "run.start", diff --git a/kernel/relayflowd/tests/crash_resume/workspace_identity.rs b/kernel/relayflowd/tests/crash_resume/workspace_identity.rs index 15c181598..c3b2f31b6 100644 --- a/kernel/relayflowd/tests/crash_resume/workspace_identity.rs +++ b/kernel/relayflowd/tests/crash_resume/workspace_identity.rs @@ -30,7 +30,7 @@ fn complete(worker: &mut ProtocolClient, dispatch: &Value) -> Value { fn workspace_aliases_are_refused_and_canonical_subtrees_serialize_over_real_sockets() { let fixture = LlmFixture::parallel("workspace-identity"); let _server = ServerGuard::start(&fixture); - let socket = fixture.data_dir.join("relayflowd.sock"); + let socket = fixture.socket(); let mut worker = ProtocolClient::connect(&socket); assert_eq!( worker.request_error_code( diff --git a/kernel/relayflowd/tests/daemon_lifecycle.rs b/kernel/relayflowd/tests/daemon_lifecycle.rs index 1b99f4b3c..8342432ae 100644 --- a/kernel/relayflowd/tests/daemon_lifecycle.rs +++ b/kernel/relayflowd/tests/daemon_lifecycle.rs @@ -90,9 +90,14 @@ fn connection_file_is_published_only_after_the_socket_is_live() { connection["protocol"].as_u64(), Some(PROTOCOL_VERSION as u64) ); + // Socket now lives outside the data dir at a derived path (#262). The + // anti-hijack invariant is unchanged: the file must name the exact same + // path the CLI would recompute from --data-dir. + let expected_socket = relayflowd::socket_path::derive_socket_path(directory.path()) + .expect("derive socket path"); assert_eq!( connection["socket_path"].as_str(), - Some(directory.path().join("relayflowd.sock").to_str().unwrap()) + expected_socket.to_str(), ); signal(&daemon, libc::SIGTERM); assert!(daemon.wait().unwrap().success()); @@ -106,7 +111,9 @@ fn clean_shutdown_removes_advertisement_and_socket() { signal(&daemon, libc::SIGINT); assert!(daemon.wait().unwrap().success()); assert!(!connection_path(directory.path()).exists()); - assert!(!directory.path().join("relayflowd.sock").exists()); + let socket = relayflowd::socket_path::derive_socket_path(directory.path()) + .expect("derive socket path"); + assert!(!socket.exists()); } #[test] @@ -191,3 +198,38 @@ fn a_sigkilled_daemons_successor_starts_cleanly() { signal(&second, libc::SIGTERM); assert!(second.wait().unwrap().success()); } + +/// #262 regression: a deep working directory used to push +/// `/relayflowd.sock` past SUN_LEN (~104 bytes on macOS). The socket +/// now lives outside the data dir at a short hashed path, so the daemon binds +/// and serves regardless of how deep the data dir is. Padding to well past +/// the limit to prove the property holds by construction, not by accident. +#[test] +fn deep_data_dir_still_binds() { + let directory = TempDir::new().unwrap(); + let mut deep = directory.path().to_path_buf(); + // Grow the data-dir path well past the ~104-byte macOS ceiling. + while deep.as_os_str().len() < 200 { + deep = deep.join("nested-directory-with-a-long-enough-name-to-add-length"); + } + assert!( + deep.as_os_str().len() > 104, + "test fixture must actually exceed SUN_LEN, got {} bytes", + deep.as_os_str().len(), + ); + std::fs::create_dir_all(&deep).unwrap(); + + let mut daemon = start(&deep); + let connection = wait_for_connection(&deep); + let socket = connection["socket_path"].as_str().unwrap(); + // The socket is what the caller must reach, so it is the length that + // matters — not the data-dir length. + assert!( + socket.len() < 104, + "derived socket path must fit SUN_LEN, got {} bytes: {socket}", + socket.len(), + ); + UnixStream::connect(socket).expect("deep-data-dir daemon must accept connections"); + signal(&daemon, libc::SIGTERM); + assert!(daemon.wait().unwrap().success()); +} diff --git a/kernel/relayflowd/tests/invalid_schema_preflight.rs b/kernel/relayflowd/tests/invalid_schema_preflight.rs index d25277502..e105a26fd 100644 --- a/kernel/relayflowd/tests/invalid_schema_preflight.rs +++ b/kernel/relayflowd/tests/invalid_schema_preflight.rs @@ -33,7 +33,8 @@ struct Refusal { fn submit_gate(schema: Value) -> Refusal { let directory = tempfile::tempdir().unwrap(); let marker = directory.path().join("command-ran"); - let socket = directory.path().join("relayflowd.sock"); + let socket = relayflowd::socket_path::derive_socket_path(directory.path()) + .expect("derive socket path"); let spec = json!({ "steps": [{ "id": "schema", diff --git a/packages/sdk/src/cli/hn-monitor.ts b/packages/sdk/src/cli/hn-monitor.ts index 0f17a88a6..7f064ee97 100644 --- a/packages/sdk/src/cli/hn-monitor.ts +++ b/packages/sdk/src/cli/hn-monitor.ts @@ -16,6 +16,7 @@ import { readFile } from 'node:fs/promises'; import { dirname, isAbsolute, resolve } from 'node:path'; import { pollHackerNewsOnce, HnTransientFetchError, type Fetcher } from '../hn-poller.js'; import { AgentWorker } from '../worker.js'; +import { socketPathFor } from '../daemon-connection.js'; import { JournalClient } from '../journal-client.js'; import type { EventSubmitResult, HelloResult, Pins } from '../protocol.js'; import type { CliIo } from '../cli.js'; @@ -181,7 +182,7 @@ async function defaultAttachWorker( */ export async function runHnMonitor(args: HnMonitorArgs, io: CliIo): Promise<0 | 1> { const pollIntervalMs = args.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS; - const socketPath = `${args.dataDir}/relayflowd.sock`; + const socketPath = socketPathFor(args.dataDir); let spec: unknown; try { diff --git a/packages/sdk/src/cli/run.ts b/packages/sdk/src/cli/run.ts index e76bb871e..c923c1bf4 100644 --- a/packages/sdk/src/cli/run.ts +++ b/packages/sdk/src/cli/run.ts @@ -1,6 +1,7 @@ import { join, resolve } from 'node:path'; import type { ProgressEvent } from '../progress.js'; import { toKernelSpec } from '../compile.js'; +import { socketPathFor } from '../daemon-connection.js'; import { ensureDaemon, type EnsureDaemonOptions } from '../daemon-lifecycle.js'; import { daemonRefusal } from './daemon-refusal.js'; import type { RunFailureKind, RunWarningKind } from '../failure-kinds.js'; @@ -464,8 +465,12 @@ function fromBase(command: RunCommand, base: CheckReport | RunReport): RunReport return 'command' in base ? base : fromCheckReport(command, base); } +// Delegates to the daemon-connection derivation so run.ts, direct-run.ts, and +// the daemon-lifecycle attach path all speak the same socket path. Kept as a +// re-export here so existing callers do not have to reach into +// daemon-connection.ts. See socketPathFor for the SUN_LEN reasoning (#262). export function socketFor(dataDir: string): string { - return join(resolve(dataDir), 'relayflowd.sock'); + return socketPathFor(dataDir); } function errorMessage(error: unknown): string { diff --git a/packages/sdk/src/cli/tick-runner.ts b/packages/sdk/src/cli/tick-runner.ts index 9063c616d..042d349b2 100644 --- a/packages/sdk/src/cli/tick-runner.ts +++ b/packages/sdk/src/cli/tick-runner.ts @@ -37,6 +37,7 @@ import { mkdir, readFile, writeFile } from 'node:fs/promises'; import { dirname, isAbsolute, join, resolve } from 'node:path'; +import { socketPathFor } from '../daemon-connection.js'; import { JournalClient } from '../journal-client.js'; import { DEFAULT_MAX_CATCH_UP, @@ -269,7 +270,7 @@ export async function runTickRunner(args: TickRunnerArgs, io: CliIo): Promise/`, ~14 chars — the shortest + // per-user private option). TMPDIR next (macOS, `/var/folders/xx/YYY/T/`, + // per-user private). Only fall back to os.tmpdir() when neither is set. + const xdg = process.env['XDG_RUNTIME_DIR']; + if (xdg && xdg.length > 0) return xdg; + const tmp = process.env['TMPDIR']; + if (tmp && tmp.length > 0) return tmp; + return tmpdir(); } export function connectionPathFor(dataDir: string): string { diff --git a/packages/sdk/src/demo-hn-monitor.ts b/packages/sdk/src/demo-hn-monitor.ts index 3cfff5587..e377e4b2e 100644 --- a/packages/sdk/src/demo-hn-monitor.ts +++ b/packages/sdk/src/demo-hn-monitor.ts @@ -1,6 +1,7 @@ import { readFile } from 'node:fs/promises'; import { dirname, join, resolve } from 'node:path'; import { fileURLToPath } from 'node:url'; +import { socketPathFor } from './daemon-connection.js'; import { pollHackerNewsOnce, type EventSink } from './hn-poller.js'; import { JournalClient } from './journal-client.js'; import type { EventSubmitResult } from './protocol.js'; @@ -8,7 +9,7 @@ import type { EventSubmitResult } from './protocol.js'; const sdkRoot = resolve(dirname(fileURLToPath(import.meta.url)), '..'); const repositoryRoot = resolve(sdkRoot, '..'); const dataDir = resolve(process.env.RELAYFLOW_DATA_DIR ?? join(repositoryRoot, '.relayflowd')); -const socketPath = join(dataDir, 'relayflowd.sock'); +const socketPath = socketPathFor(dataDir); const specPath = join(repositoryRoot, 'testdata', 'hn-monitor.spec.canonical.json'); interface Submission { diff --git a/packages/sdk/tests/cli.test.ts b/packages/sdk/tests/cli.test.ts index b55d6b5c7..52ba43cad 100644 --- a/packages/sdk/tests/cli.test.ts +++ b/packages/sdk/tests/cli.test.ts @@ -8,6 +8,7 @@ import { parse as parseYaml, stringify as stringifyYaml } from 'yaml'; import { afterEach, describe, expect, it } from 'vitest'; import { runCli, type CheckReport, type CliIo } from '../src/cli.js'; import { checkFlow } from '../src/cli/check.js'; +import { socketPathFor } from '../src/daemon-connection.js'; import { runFlow } from '../src/cli/run.js'; import { runAgentCli } from '../src/worker-cli.js'; import { @@ -139,7 +140,7 @@ const LADDER_FAULTS = [ ] as const satisfies ReadonlyArray) => void]>; async function startCliLoopback(dataDir: string, handlers: LoopbackHandlers): Promise { - const server = startLoopback(join(dataDir, 'relayflowd.sock'), handlers); + const server = startLoopback(socketPathFor(dataDir), handlers); loopbackServers.push(server); if (!server.listening) await once(server, 'listening'); } diff --git a/packages/sdk/tests/daemon-lifecycle-live.test.ts b/packages/sdk/tests/daemon-lifecycle-live.test.ts index 443c6fbca..1ede5335a 100644 --- a/packages/sdk/tests/daemon-lifecycle-live.test.ts +++ b/packages/sdk/tests/daemon-lifecycle-live.test.ts @@ -24,6 +24,7 @@ import { dirname, join } from 'node:path'; import { fileURLToPath } from 'node:url'; import { afterEach, beforeAll, describe, expect, it } from 'vitest'; import type { DaemonConnection } from '../src/daemon-lifecycle.js'; +import { socketPathFor } from '../src/daemon-connection.js'; const HERE = dirname(fileURLToPath(import.meta.url)); const ROOT = join(HERE, '..', '..', '..'); @@ -147,7 +148,7 @@ describe('flows run against a data dir with no daemon (§6 test 7)', () => { // daemon it started is still serving. const connection = connectionFile(dataDir); startedPids.push(connection.pid); - expect(connection.socket_path).toBe(join(dataDir, 'relayflowd.sock')); + expect(connection.socket_path).toBe(socketPathFor(dataDir)); expect(isAlive(connection.pid)).toBe(true); }); @@ -187,7 +188,7 @@ describe('flows run against a data dir with no daemon (§6 test 7)', () => { expect(result.status, result.stderr).toBe(0); expect(result.stdout).toContain('completionReason: success'); expect(existsSync(join(dataDir, 'connection.json'))).toBe(false); - expect(existsSync(join(dataDir, 'relayflowd.sock'))).toBe(true); + expect(existsSync(socketPathFor(dataDir))).toBe(true); startedPids.push(lockHolder(dataDir)); }); diff --git a/packages/sdk/tests/direct-input.test.ts b/packages/sdk/tests/direct-input.test.ts index 0ccc0475e..ade726805 100644 --- a/packages/sdk/tests/direct-input.test.ts +++ b/packages/sdk/tests/direct-input.test.ts @@ -14,6 +14,7 @@ import { dirname, join, resolve } from 'node:path'; import { spawn, spawnSync, type ChildProcess } from 'node:child_process'; import { fileURLToPath } from 'node:url'; import { afterEach, beforeAll, describe, expect, it } from 'vitest'; +import { socketPathFor } from '../src/daemon-connection.js'; const ROOT = join(dirname(fileURLToPath(import.meta.url)), '..', '..', '..'); const BUILT_CLI = join(ROOT, 'packages', 'sdk', 'dist', 'cli.js'); @@ -154,7 +155,7 @@ async function startDaemon(dataDir: string): Promise { stdio: ['ignore', 'pipe', 'pipe'], }); daemons.push(daemon); - const socket = join(dataDir, 'relayflowd.sock'); + const socket = socketPathFor(dataDir); for (let attempt = 0; attempt < 100; attempt += 1) { if (existsSync(socket)) return; if (daemon.exitCode !== null) throw new Error(`relayflowd exited ${daemon.exitCode}`); diff --git a/packages/sdk/tests/fixtures/stub-relayflowd.mjs b/packages/sdk/tests/fixtures/stub-relayflowd.mjs index def83e100..0191a6b14 100644 --- a/packages/sdk/tests/fixtures/stub-relayflowd.mjs +++ b/packages/sdk/tests/fixtures/stub-relayflowd.mjs @@ -18,6 +18,7 @@ // `flock` guarantee itself belongs to the kernel implementation and its // `cargo test` cases (DAEMON-LIFECYCLE.md §6 tests 4 and 5). +import { createHash } from 'node:crypto'; import { createServer } from 'node:net'; import { appendFileSync, @@ -31,8 +32,24 @@ import { writeFileSync, writeSync, } from 'node:fs'; +import { tmpdir } from 'node:os'; import { join, resolve } from 'node:path'; +// Must stay in lockstep with `packages/sdk/src/daemon-connection.ts::socketPathFor` +// and `kernel/relayflowd/src/socket_path.rs::derive_socket_path` (see #262). The +// CLI derives the same path from the same input; anti-hijack (DAEMON-LIFECYCLE +// §2 step 3) breaks the moment these diverge, so update all three together. +function deriveSocketPath(dataDir) { + const absolute = resolve(dataDir); + const hash = createHash('sha256').update(absolute).digest('hex').slice(0, 12); + const runtime = (process.env.XDG_RUNTIME_DIR && process.env.XDG_RUNTIME_DIR.length > 0) + ? process.env.XDG_RUNTIME_DIR + : ((process.env.TMPDIR && process.env.TMPDIR.length > 0) + ? process.env.TMPDIR + : tmpdir()); + return join(runtime, `relayflowd-${hash}.sock`); +} + const EXIT_ALREADY_SERVING = 3; const VERSION = '0.0.0-stub'; @@ -43,7 +60,7 @@ if (dataDirIndex === -1 || argv[dataDirIndex + 1] === undefined) { fail('stub relayflowd: expected --data-dir '); } const dataDir = resolve(argv[dataDirIndex + 1]); -const socketPath = join(dataDir, 'relayflowd.sock'); +const socketPath = deriveSocketPath(dataDir); const connectionPath = join(dataDir, 'connection.json'); const lockPath = join(dataDir, 'relayflowd.lock'); diff --git a/packages/sdk/tests/live-kernel.test.ts b/packages/sdk/tests/live-kernel.test.ts index 54aa37b00..ab520d904 100644 --- a/packages/sdk/tests/live-kernel.test.ts +++ b/packages/sdk/tests/live-kernel.test.ts @@ -19,6 +19,7 @@ import { afterEach, beforeAll, describe, expect, it } from 'vitest'; import { flow } from '@relayflows/surface'; import { compileYaml, toKernelSpec } from '../src/compile.js'; import { checkFlow } from '../src/cli/check.js'; +import { socketPathFor } from '../src/daemon-connection.js'; import { JournalClient } from '../src/journal-client.js'; import type { StepDispatchEvent } from '../src/protocol.js'; import { AgentWorker } from '../src/worker.js'; @@ -1322,7 +1323,7 @@ steps: // this case pins -- refused before any journal write, naming the socket -- // is unchanged, and `runArtifacts(absentDir)` below still proves it. const absentDir = temporaryDirectory('flows-live-absent-'); - const absentSocket = join(absentDir, 'relayflowd.sock'); + const absentSocket = socketPathFor(absentDir); const unreachable = invokeCli([ 'run', '--no-spawn', '--data-dir', absentDir, join(TESTDATA, 'hello-deterministic.flow.yaml'), ]); @@ -1780,7 +1781,7 @@ async function startDaemon(dataDir: string): Promise { daemons.push(daemon); const stderr: Buffer[] = []; daemon.stderr?.on('data', (chunk: Buffer) => stderr.push(chunk)); - const socket = join(dataDir, 'relayflowd.sock'); + const socket = socketPathFor(dataDir); const deadline = Date.now() + 5_000; while (Date.now() < deadline) { if (existsSync(socket) && lstatSync(socket).isSocket()) return daemon; @@ -1800,7 +1801,7 @@ async function stopDaemon(daemon: ChildProcess, signal: NodeJS.Signals = 'SIGTER } async function connectClient(dataDir: string): Promise { - const client = new JournalClient(join(dataDir, 'relayflowd.sock'), { requestTimeoutMs: 5_000 }); + const client = new JournalClient(socketPathFor(dataDir), { requestTimeoutMs: 5_000 }); clients.push(client); await client.connect(); return client;