use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use super::state::HarnessState;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PortHolder {
pub pid: u32,
pub cwd: PathBuf,
pub env_path: Option<PathBuf>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PortOwnership {
Free,
Ours(PortHolder),
Foreign(PortHolder),
}
pub fn classify_holder(holder: PortHolder, state: &HarnessState) -> PortOwnership {
let pid_recorded = state.instances.iter().any(|inst| inst.pid == Some(holder.pid));
let env_under_known_instance = holder.env_path.as_deref().is_some_and(|env_path| {
state.instances.iter().any(|inst| {
instance_root(&inst.session_path)
.map(|root| env_path.starts_with(&root))
.unwrap_or(false)
})
});
if pid_recorded || env_under_known_instance {
PortOwnership::Ours(holder)
} else {
PortOwnership::Foreign(holder)
}
}
fn instance_root(session_path: &Path) -> Option<PathBuf> {
session_path.parent().and_then(Path::parent).map(Path::to_path_buf)
}
pub fn classify_port_holder(
holder: PortHolder,
expected_env_path: Option<&Path>,
state: Option<&HarnessState>,
) -> PortOwnership {
let owns_by_env = matches!(
(holder.env_path.as_deref(), expected_env_path),
(Some(actual), Some(expected)) if actual == expected
);
let owns_by_state = state.is_some_and(|s| {
matches!(classify_holder(holder.clone(), s), PortOwnership::Ours(_))
});
if owns_by_env || owns_by_state {
PortOwnership::Ours(holder)
} else {
PortOwnership::Foreign(holder)
}
}
pub fn classify_port(
port: u16,
expected_env_path: Option<&Path>,
state: Option<&HarnessState>,
) -> PortOwnership {
match discover_port_holder(port) {
None => PortOwnership::Free,
Some(holder) => classify_port_holder(holder, expected_env_path, state),
}
}
pub fn lane_offset(instance: &str, client_node: bool) -> u16 {
if !client_node {
return 0;
}
match instance {
"alice" => crate::commands::dev::host::monorepo::CLIENT_NODE_PORT_OFFSET,
"bob" => crate::commands::dev::host::monorepo::CLIENT_NODE_PORT_OFFSET + 100,
other => 500 + (stable_hash(other) % 400),
}
}
fn stable_hash(s: &str) -> u16 {
use std::hash::{Hash, Hasher};
let mut h = std::collections::hash_map::DefaultHasher::new();
s.hash(&mut h);
(h.finish() % (u16::MAX as u64)) as u16
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct PortSet {
pub http: u16,
pub https: u16,
pub p2p: u16,
pub ui: u16,
}
impl PortSet {
pub fn base_for(
instance: &str,
http: u16,
https: u16,
p2p: u16,
ui: u16,
) -> Self {
let offset = lane_offset(instance, true);
Self { http: http + offset, https: https + offset, p2p: p2p + offset, ui: ui + offset }
}
fn shifted(self, delta: u16) -> Self {
Self {
http: self.http + delta,
https: self.https + delta,
p2p: self.p2p + delta,
ui: self.ui + delta,
}
}
fn fields(self) -> [(&'static str, u16); 4] {
[("http", self.http), ("https", self.https), ("p2p", self.p2p), ("ui", self.ui)]
}
}
const LANE_SLOT_STRIDE: u16 = 50;
const MAX_LANE_ATTEMPTS: u16 = 5;
pub fn allocate_lane_ports(
base: PortSet,
expected_env_path: Option<&Path>,
state: Option<&HarnessState>,
) -> anyhow::Result<PortSet> {
allocate_lane_ports_with(base, |port| classify_port(port, expected_env_path, state))
}
fn allocate_lane_ports_with(
base: PortSet,
mut classify: impl FnMut(u16) -> PortOwnership,
) -> anyhow::Result<PortSet> {
let mut tried = Vec::new();
for attempt in 0..MAX_LANE_ATTEMPTS {
let candidate = base.shifted(attempt * LANE_SLOT_STRIDE);
let mut blockers = Vec::new();
for (label, port) in candidate.fields() {
if let PortOwnership::Foreign(h) = classify(port) {
blockers.push(format!(
"{label} {port} held by pid {} (cwd {}, env {})",
h.pid,
h.cwd.display(),
h.env_path
.as_deref()
.map(|p| p.display().to_string())
.unwrap_or_else(|| "<unknown>".into()),
));
}
}
if blockers.is_empty() {
return Ok(candidate);
}
tried.push(format!("slot {attempt} (http {}): {}", candidate.http, blockers.join("; ")));
}
anyhow::bail!(
"no free client-node lane found after {MAX_LANE_ATTEMPTS} candidate slot(s) starting \
from http {}:\n{}\n\
Every candidate was held by a process this invocation cannot attribute to itself. If \
one of those is actually yours, stop it (`node-app harness down` for a stale harness) \
and re-run; otherwise another session is genuinely using this box's client-node lanes \
and there is no free slot left to allocate.",
base.http,
tried.join("\n"),
);
}
pub fn discover_port_holder(port: u16) -> Option<PortHolder> {
let pid = listening_pid(port)?;
Some(describe_pid(pid))
}
fn listening_pid(port: u16) -> Option<u32> {
let arg = format!("-iTCP:{port}");
let output = Command::new("lsof")
.args(["-nP", "-sTCP:LISTEN", &arg])
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.output()
.ok()?;
if !output.status.success() {
return None;
}
parse_listening_pid(&String::from_utf8_lossy(&output.stdout))
}
fn parse_listening_pid(listing: &str) -> Option<u32> {
listing
.lines()
.skip(1) .find_map(|line| line.split_whitespace().nth(1).and_then(|p| p.parse::<u32>().ok()))
}
pub fn describe_pid(pid: u32) -> PortHolder {
PortHolder { pid, cwd: pid_cwd(pid).unwrap_or_default(), env_path: pid_env_arg(pid) }
}
fn pid_cwd(pid: u32) -> Option<PathBuf> {
let output = Command::new("lsof")
.args(["-a", "-p", &pid.to_string(), "-d", "cwd", "-Fn"])
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.output()
.ok()?;
if !output.status.success() {
return None;
}
parse_lsof_cwd(&String::from_utf8_lossy(&output.stdout))
}
fn parse_lsof_cwd(output: &str) -> Option<PathBuf> {
output.lines().find_map(|l| l.strip_prefix('n')).filter(|p| !p.is_empty()).map(PathBuf::from)
}
fn pid_env_arg(pid: u32) -> Option<PathBuf> {
let output = Command::new("ps")
.args(["-o", "command=", "-p", &pid.to_string()])
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.output()
.ok()?;
if !output.status.success() {
return None;
}
parse_env_arg(&String::from_utf8_lossy(&output.stdout))
}
fn parse_env_arg(command: &str) -> Option<PathBuf> {
let rest = command.split("--env ").nth(1)?;
let env_path = rest.split_whitespace().next()?;
if env_path.is_empty() {
return None;
}
Some(PathBuf::from(env_path))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::commands::harness::state::InstanceState;
#[test]
fn a_holder_absent_from_harness_state_is_foreign() {
let state = HarnessState::default();
let holder = PortHolder { pid: 4242, cwd: "/tmp/other".into(), env_path: None };
assert!(matches!(classify_holder(holder, &state), PortOwnership::Foreign(_)));
}
#[test]
fn a_holder_recorded_in_harness_state_is_ours() {
let mut state = HarnessState::default();
state.instances.push(InstanceState { pid: Some(4242), ..InstanceState::default() });
let holder = PortHolder { pid: 4242, cwd: "/tmp/mine".into(), env_path: None };
assert!(matches!(classify_holder(holder, &state), PortOwnership::Ours(_)));
}
#[test]
fn env_path_under_a_known_instance_root_is_ours_even_with_a_different_pid() {
let mut state = HarnessState::default();
state.instances.push(InstanceState {
session_path: "/cache/monorepo-abc/dev-apps/alice-agent-session.json".into(),
..InstanceState::default()
});
let holder = PortHolder {
pid: 9999, cwd: PathBuf::new(),
env_path: Some("/cache/monorepo-abc/daemon.env".into()),
};
assert!(matches!(classify_holder(holder, &state), PortOwnership::Ours(_)));
}
#[test]
fn a_holder_with_unparseable_attribution_fields_is_foreign() {
let state = HarnessState::default();
let holder = PortHolder { pid: 777, cwd: PathBuf::new(), env_path: None };
assert!(matches!(classify_holder(holder, &state), PortOwnership::Foreign(_)));
}
#[test]
fn matching_env_path_alone_is_ours_even_without_a_state_file() {
let expected = PathBuf::from("/cache/node-app/monorepo-abc/daemon.env");
let holder = PortHolder { pid: 555, cwd: PathBuf::new(), env_path: Some(expected.clone()) };
assert!(matches!(
classify_port_holder(holder, Some(&expected), None),
PortOwnership::Ours(_)
));
}
#[test]
fn a_different_checkouts_env_path_is_foreign_even_with_no_state() {
let expected = PathBuf::from("/cache/node-app/monorepo-AAA/daemon.env");
let foreign_env = PathBuf::from("/cache/node-app/monorepo-BBB/daemon.env");
let holder = PortHolder { pid: 555, cwd: PathBuf::new(), env_path: Some(foreign_env) };
assert!(matches!(
classify_port_holder(holder, Some(&expected), None),
PortOwnership::Foreign(_)
));
}
#[test]
fn no_expected_path_and_no_state_is_foreign_never_ours() {
let holder = PortHolder { pid: 1, cwd: PathBuf::new(), env_path: None };
assert!(matches!(classify_port_holder(holder, None, None), PortOwnership::Foreign(_)));
}
#[test]
fn lane_offsets_are_distinct_per_instance() {
assert_ne!(lane_offset("alice", true), lane_offset("bob", true));
assert_eq!(lane_offset("alice", false), 0);
}
#[test]
fn plain_lane_is_always_offset_zero_regardless_of_instance() {
assert_eq!(lane_offset("alice", false), 0);
assert_eq!(lane_offset("bob", false), 0);
assert_eq!(lane_offset("someone-else", false), 0);
}
#[test]
fn alice_client_node_offset_matches_the_existing_documented_default() {
assert_eq!(lane_offset("alice", true), 300);
}
#[test]
fn parse_listening_pid_reads_the_pid_column() {
let listing = "COMMAND PID USER FD TYPE DEVICE SIZE/OFF NODE NAME\n\
node-serv 4242 vulam 6u IPv4 0x123 0t0 TCP *:3001 (LISTEN)\n";
assert_eq!(parse_listening_pid(listing), Some(4242));
}
#[test]
fn parse_listening_pid_is_none_on_garbage() {
assert_eq!(parse_listening_pid("not lsof output at all\n***\ngarbage"), None);
assert_eq!(parse_listening_pid(""), None);
}
#[test]
fn parse_lsof_cwd_reads_the_n_field() {
let out = "p4242\nfcwd\nn/Users/vulam/checkout\n";
assert_eq!(parse_lsof_cwd(out), Some(PathBuf::from("/Users/vulam/checkout")));
}
#[test]
fn parse_lsof_cwd_is_none_on_garbage() {
assert_eq!(parse_lsof_cwd("garbage output\nsomething else entirely"), None);
}
#[test]
fn parse_env_arg_extracts_the_env_flag_value() {
let cmd = "/path/to/node-server --env /cache/node-app/monorepo-abc/daemon.env";
assert_eq!(
parse_env_arg(cmd),
Some(PathBuf::from("/cache/node-app/monorepo-abc/daemon.env"))
);
}
#[test]
fn parse_env_arg_is_none_without_the_flag() {
assert_eq!(parse_env_arg("/usr/bin/some-other-process --foo bar"), None);
assert_eq!(parse_env_arg(""), None);
}
#[test]
fn classify_port_is_free_when_nothing_is_listening() {
assert_eq!(classify_port(59_999, None, None), PortOwnership::Free);
}
#[test]
#[cfg(unix)]
fn a_real_listening_socket_with_no_matching_env_is_foreign() {
use std::net::TcpListener;
let Ok(listener) = TcpListener::bind("127.0.0.1:0") else {
return; };
let port = listener.local_addr().unwrap().port();
let holder = discover_port_holder(port);
drop(listener);
let Some(holder) = holder else {
return;
};
let expected = PathBuf::from("/definitely/not/our/daemon.env");
assert!(matches!(
classify_port_holder(holder, Some(&expected), None),
PortOwnership::Foreign(_)
));
}
fn sample_base() -> PortSet {
PortSet { http: 3301, https: 4731, p2p: 10_035, ui: 5_473 }
}
fn free_holder() -> PortHolder {
PortHolder { pid: 1, cwd: PathBuf::new(), env_path: None }
}
#[test]
fn allocate_lane_ports_takes_the_documented_base_when_everything_is_free() {
let base = sample_base();
let result = allocate_lane_ports_with(base, |_port| PortOwnership::Free).unwrap();
assert_eq!(result, base, "the common case — one harness on a box — must keep the default");
}
#[test]
fn allocate_lane_ports_takes_the_documented_base_when_it_is_ours() {
let base = sample_base();
let holder = free_holder();
let result =
allocate_lane_ports_with(base, move |_port| PortOwnership::Ours(holder.clone()))
.unwrap();
assert_eq!(result, base);
}
#[test]
fn allocate_lane_ports_skips_a_foreign_base_and_resolves_to_the_next_slot() {
let base = sample_base();
let foreign = PortHolder {
pid: 999,
cwd: PathBuf::from("/tmp/other-checkout"),
env_path: Some(PathBuf::from("/cache/node-app/monorepo-OTHER/daemon.env")),
};
let result = allocate_lane_ports_with(base, move |port| {
if port == base.http {
PortOwnership::Foreign(foreign.clone())
} else {
PortOwnership::Free
}
})
.unwrap();
assert_ne!(result, base, "must not choose a slot with a foreign holder");
assert_eq!(result, base.shifted(LANE_SLOT_STRIDE));
}
#[test]
fn allocate_lane_ports_rejects_the_whole_slot_when_only_one_port_is_foreign() {
let base = sample_base();
let foreign = PortHolder { pid: 999, cwd: PathBuf::new(), env_path: Some("/other".into()) };
let result = allocate_lane_ports_with(base, move |port| {
if port == base.https {
PortOwnership::Foreign(foreign.clone())
} else {
PortOwnership::Free
}
})
.unwrap();
assert_eq!(result, base.shifted(LANE_SLOT_STRIDE));
}
#[test]
fn allocate_lane_ports_fails_after_exhausting_attempts_and_lists_every_candidate_and_holder() {
let base = sample_base();
let foreign =
PortHolder { pid: 999, cwd: "/tmp/blocker".into(), env_path: Some("/other".into()) };
let err =
allocate_lane_ports_with(base, move |_port| PortOwnership::Foreign(foreign.clone()))
.unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("999"), "must name the blocking pid: {msg}");
assert!(msg.contains("/tmp/blocker"), "must name the blocking cwd: {msg}");
for attempt in 0..MAX_LANE_ATTEMPTS {
let candidate_http = base.http + attempt * LANE_SLOT_STRIDE;
assert!(
msg.contains(&candidate_http.to_string()),
"must list every candidate tried, missing slot {attempt} ({candidate_http}): {msg}"
);
}
}
#[test]
fn allocate_lane_ports_never_terminates_anything() {
let base = sample_base();
let foreign = free_holder();
let _ = allocate_lane_ports_with(base, move |_port| PortOwnership::Foreign(foreign.clone()));
}
#[test]
#[cfg(unix)]
fn allocate_lane_ports_resolves_a_second_lane_around_a_real_foreign_holder() {
use std::net::TcpListener;
let Ok(l_http) = TcpListener::bind("127.0.0.1:0") else { return };
let http_port = l_http.local_addr().unwrap().port();
if discover_port_holder(http_port).is_none() {
return;
}
let (Ok(l_https), Ok(l_p2p), Ok(l_ui)) = (
TcpListener::bind("127.0.0.1:0"),
TcpListener::bind("127.0.0.1:0"),
TcpListener::bind("127.0.0.1:0"),
) else {
return;
};
let base = PortSet {
http: http_port,
https: l_https.local_addr().unwrap().port(),
p2p: l_p2p.local_addr().unwrap().port(),
ui: l_ui.local_addr().unwrap().port(),
};
let expected = PathBuf::from("/definitely/not/our/daemon.env");
let result = allocate_lane_ports(base, Some(&expected), None);
drop((l_http, l_https, l_p2p, l_ui));
let result = result.expect("a free slot exists well within MAX_LANE_ATTEMPTS");
assert_ne!(result, base, "must not choose the slot this test process itself occupies");
assert_eq!(result, base.shifted(LANE_SLOT_STRIDE));
}
#[test]
fn port_set_base_for_matches_lane_offset() {
let base = PortSet::base_for("alice", 3001, 4431, 9735, 5173);
assert_eq!(base.http, 3001 + lane_offset("alice", true));
assert_eq!(base.https, 4431 + lane_offset("alice", true));
assert_eq!(base.p2p, 9735 + lane_offset("alice", true));
assert_eq!(base.ui, 5173 + lane_offset("alice", true));
}
}