use std::{
net::SocketAddr,
path::PathBuf,
time::{Duration, Instant},
};
use clap::Args;
use eyre::Context;
use crate::{
command::{Executable, default_tracing, up},
common::connect_to_coordinator,
};
use super::config::{ClusterConfig, ZenohMesh};
use super::{
format_daemon_port_arg, format_labels_arg, format_zenoh_peer_arg, query_connected_daemons,
run_ssh, ssh_target,
};
#[derive(Debug, Args)]
#[clap(verbatim_doc_comment)]
pub struct Up {
#[clap(value_name = "PATH", value_hint = clap::ValueHint::FilePath)]
config: PathBuf,
}
impl Executable for Up {
fn execute(self) -> eyre::Result<()> {
default_tracing()?;
let config = ClusterConfig::load(&self.config)?;
let coordinator_addr: SocketAddr =
(config.coordinator.addr, config.coordinator.port).into();
let session = match connect_to_coordinator(coordinator_addr) {
Ok(s) => {
println!("Coordinator already running at {coordinator_addr}");
s
}
Err(_) => {
let startup = up::spawn_coordinator(up::CoordinatorSpawn {
interface: Some(std::net::Ipv4Addr::UNSPECIFIED.into()),
port: Some(config.coordinator.port),
auth: false,
})
.wrap_err("failed to start dora coordinator")?;
up::wait_for_coordinator_start(coordinator_addr, config.coordinator.port, startup)?
}
};
let zenoh_peer_arg = format_zenoh_peer_arg(config.zenoh_peer.as_deref());
let zenoh_mesh_args = match config.zenoh_mesh_args() {
ZenohMesh::Derived(args) => Some(args),
ZenohMesh::NotNeeded => None,
ZenohMesh::Unavailable(reason) => {
eprintln!(
"WARNING: {reason}, so the daemons are left to discover each other \
by multicast. On a network without multicast — a mesh VPN carries \
none — they will not find each other. Fix the field named above, \
or configure a shared `zenoh_peer` rendezvous."
);
None
}
};
let mut ssh_failures: Vec<(String, String)> = Vec::new();
for machine in &config.machines {
let target = ssh_target(machine);
let labels_arg = format_labels_arg(&machine.labels);
let daemon_port_arg = format_daemon_port_arg(machine.daemon_port);
let mesh_arg = zenoh_mesh_args
.as_ref()
.and_then(|args| args.get(machine.id.as_str()))
.map(String::as_str)
.unwrap_or_default();
let remote_cmd = format!(
"nohup dora daemon --machine-id {id} --coordinator-addr {addr} --coordinator-port {port}{daemon_port_arg}{zenoh_peer_arg}{mesh_arg}{labels} --quiet > /tmp/dora-daemon-{id}.log 2>&1 &",
id = machine.id,
addr = config.coordinator.addr,
port = config.coordinator.port,
labels = labels_arg,
);
println!("Starting daemon on {} ({})", machine.id, target);
match run_ssh(&target, machine.port, &remote_cmd) {
Ok(true) => {}
Ok(false) => {
let msg = "ssh command failed".to_string();
eprintln!(
" WARNING: failed to start daemon on `{}`: {msg}",
machine.id
);
ssh_failures.push((machine.id.clone(), msg));
}
Err(err) => {
let msg = format!("{err}");
eprintln!(
" WARNING: failed to start daemon on `{}`: {msg}",
machine.id
);
ssh_failures.push((machine.id.clone(), msg));
}
}
}
let expected: Vec<&str> = config
.machines
.iter()
.filter(|m| !ssh_failures.iter().any(|(id, _)| id == &m.id))
.map(|m| m.id.as_str())
.collect();
let mut missing_daemons: Vec<String> = Vec::new();
if !expected.is_empty() {
println!("Waiting for {} daemon(s) to connect...", expected.len());
let deadline = Instant::now() + Duration::from_secs(30);
loop {
let connected = query_connected_daemons(&session)?;
let all_present = expected.iter().all(|machine_id| {
connected
.iter()
.any(|d| d.daemon_id.matches_machine_id(machine_id))
});
if all_present {
break;
}
if Instant::now() >= deadline {
missing_daemons = expected
.iter()
.copied()
.filter(|machine_id| {
!connected
.iter()
.any(|d| d.daemon_id.matches_machine_id(machine_id))
})
.map(String::from)
.collect();
eprintln!(
"WARNING: timed out waiting for daemon(s): {}",
missing_daemons.join(", ")
);
report_missing_daemon_logs(&config, &missing_daemons);
break;
}
std::thread::sleep(Duration::from_millis(500));
}
}
let ok_count = config.machines.len() - ssh_failures.len() - missing_daemons.len();
if ssh_failures.is_empty() && missing_daemons.is_empty() {
println!(
"Cluster is up: coordinator + {} daemon(s)",
config.machines.len()
);
Ok(())
} else {
println!(
"Cluster partially up: coordinator + {ok_count}/{} daemon(s)",
config.machines.len()
);
for (id, reason) in &ssh_failures {
eprintln!(" {id}: {reason}");
}
eyre::bail!(
"cluster up incomplete: {} ssh failure(s), {} daemon(s) did not register",
ssh_failures.len(),
missing_daemons.len()
)
}
}
}
fn report_missing_daemon_logs(config: &ClusterConfig, missing: &[String]) {
for machine_id in missing {
let Some(machine) = config.machines.iter().find(|m| &m.id == machine_id) else {
continue;
};
let cmd = format!("tail -n 20 /tmp/dora-daemon-{machine_id}.log 2>/dev/null");
eprintln!(" --- last log lines from `{machine_id}` ---");
if run_ssh(&ssh_target(machine), machine.port, &cmd).is_err() {
eprintln!(" (could not read the log — ssh to `{machine_id}` failed)");
}
}
}