use super::system::status::daemon_running;
use super::{Executable, default_tracing};
use crate::{
LOCALHOST,
common::{connect_to_coordinator, connect_with_retry, send_control_request},
};
use dora_core::topics::DORA_COORDINATOR_PORT_WS_DEFAULT;
use dora_message::{
cli_to_coordinator::ControlRequest,
coordinator_to_cli::{ControlRequestReply, DataflowIdAndName},
};
use eyre::{Context, ContextCompat, bail};
use std::io::{Read, Seek};
use std::path::PathBuf;
use std::process::Stdio;
use std::{fs, net::SocketAddr, path::Path, process::Command, time::Duration};
#[derive(Debug, clap::Args)]
pub struct Up {
#[clap(long, hide = true, value_name = "PATH", value_hint = clap::ValueHint::FilePath)]
config: Option<PathBuf>,
#[clap(long)]
auth: bool,
#[clap(long, value_name = "IP", env = "DORA_COORDINATOR_INTERFACE")]
interface: Option<std::net::IpAddr>,
#[clap(long)]
recreate_store: bool,
}
impl Executable for Up {
fn execute(self) -> eyre::Result<()> {
default_tracing()?;
up(
self.config.as_deref(),
self.auth,
self.recreate_store,
self.interface,
)
}
}
#[derive(Debug, Default, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct UpConfig {}
pub(crate) fn up(
config_path: Option<&Path>,
auth: bool,
recreate_store: bool,
interface: Option<std::net::IpAddr>,
) -> eyre::Result<()> {
let UpConfig {} = parse_dora_config(config_path)?;
let env_addr: Option<std::net::IpAddr> = coordinator_env_opt("DORA_COORDINATOR_ADDR")?;
let addr = connect_addr_for(env_addr, interface)?;
let port: u16 =
coordinator_env_value("DORA_COORDINATOR_PORT", DORA_COORDINATOR_PORT_WS_DEFAULT)?;
let coordinator_addr = (addr, port).into();
let spawn = CoordinatorSpawn {
interface,
port: Some(port),
auth,
};
if recreate_store
&& let Some(env_addr) = env_addr
&& !env_addr.is_loopback()
{
bail!(
"--recreate-store only applies to the local default coordinator store, \
but DORA_COORDINATOR_ADDR is set to the non-loopback address {env_addr}\n\n \
hint: unset DORA_COORDINATOR_ADDR, or run this command on the machine \
that hosts the coordinator store"
);
}
let mut _recreation_lock = None;
let session = match connect_to_coordinator(coordinator_addr) {
Ok(session) => attach_to_running_coordinator(session, coordinator_addr, recreate_store),
Err(_) if recreate_store => {
_recreation_lock = Some(lock_default_coordinator_store_recreation()?);
match connect_to_coordinator(coordinator_addr) {
Ok(session) => {
attach_to_running_coordinator(session, coordinator_addr, recreate_store)
}
Err(err) => {
ensure_no_listener_before_archive(coordinator_addr, err)?;
archive_default_coordinator_store()?;
start_and_wait_for_coordinator(coordinator_addr, port, spawn)?
}
}
}
Err(_) => {
ensure_no_coordinator_elsewhere_on_port(coordinator_addr, interface)?;
start_and_wait_for_coordinator(coordinator_addr, port, spawn)?
}
};
if !daemon_running(&session)? {
let daemon_coordinator_addr =
Some(addr).filter(|i| !i.is_loopback() && !i.is_unspecified());
if interface.is_some_and(|i| i.is_unspecified()) {
println!(
"note: the local daemon's zenoh listener stays on loopback with a wildcard \
--interface; pass this machine's address (e.g. --interface 192.168.1.10) \
so daemons on other machines can dial its nodes directly"
);
}
start_daemon(daemon_coordinator_addr, port).wrap_err("failed to start dora-daemon")?;
let mut i = 0;
const WAIT_S: f32 = 0.1;
loop {
if daemon_running(&session)? {
break;
}
i += 1;
if i > 20 {
eyre::bail!("daemon not connected after {}s", WAIT_S * i as f32);
}
std::thread::sleep(Duration::from_secs_f32(WAIT_S));
}
} else {
println!("daemon already running");
}
Ok(())
}
fn coordinator_env_value<T>(var_name: &str, default: T) -> eyre::Result<T>
where
T: std::str::FromStr,
T::Err: std::fmt::Display,
{
parse_coordinator_env(var_name, std::env::var(var_name).ok().as_deref(), default)
}
fn connect_addr_for(
env_addr: Option<std::net::IpAddr>,
interface: Option<std::net::IpAddr>,
) -> eyre::Result<std::net::IpAddr> {
let bound_concretely = interface.filter(|i| !i.is_unspecified());
match (env_addr, bound_concretely) {
(Some(env), Some(iface)) if env != iface => bail!(
"DORA_COORDINATOR_ADDR is {env} but --interface is {iface}, so the \
coordinator would bind {iface} and answer nowhere else while every \
command looked for it on {env}\n\n \
hint: drop one of them, or set DORA_COORDINATOR_ADDR={iface}"
),
(Some(env), _) => Ok(env),
(None, Some(iface)) => Ok(iface),
(None, None) => Ok(LOCALHOST),
}
}
fn coordinator_env_opt<T>(var_name: &str) -> eyre::Result<Option<T>>
where
T: std::str::FromStr,
T::Err: std::fmt::Display,
{
parse_coordinator_env_opt(var_name, std::env::var(var_name).ok().as_deref())
}
fn parse_coordinator_env<T>(var_name: &str, raw: Option<&str>, default: T) -> eyre::Result<T>
where
T: std::str::FromStr,
T::Err: std::fmt::Display,
{
Ok(parse_coordinator_env_opt(var_name, raw)?.unwrap_or(default))
}
fn parse_coordinator_env_opt<T>(var_name: &str, raw: Option<&str>) -> eyre::Result<Option<T>>
where
T: std::str::FromStr,
T::Err: std::fmt::Display,
{
match raw {
Some(s) => s
.parse()
.map(Some)
.map_err(|e| eyre::eyre!("invalid {var_name}: {s:?} ({e})")),
None => Ok(None),
}
}
fn attach_to_running_coordinator(
session: crate::ws_client::WsSession,
coordinator_addr: SocketAddr,
recreate_store: bool,
) -> crate::ws_client::WsSession {
println!("coordinator already running on {coordinator_addr}");
if recreate_store {
println!(
"--recreate-store skipped: the running coordinator holds the store; \
run `dora down` first, then rerun `dora up --recreate-store`"
);
}
session
}
fn ensure_no_listener_before_archive(
coordinator_addr: SocketAddr,
connect_err: eyre::Report,
) -> eyre::Result<()> {
let loopback = SocketAddr::from((LOCALHOST, coordinator_addr.port()));
let mut probes = vec![coordinator_addr];
if coordinator_addr != loopback {
probes.push(loopback);
}
for addr in probes {
if std::net::TcpStream::connect_timeout(&addr, Duration::from_secs(1)).is_ok() {
return Err(connect_err.wrap_err(format!(
"refusing to archive the coordinator store: a process is listening on \
{addr} but the connection failed\n\n \
hint: resolve the connection error below (e.g. an auth token mismatch), \
or stop the coordinator with `dora down` before recreating the store"
)));
}
}
Ok(())
}
fn ensure_no_coordinator_elsewhere_on_port(
coordinator_addr: SocketAddr,
interface: Option<std::net::IpAddr>,
) -> eyre::Result<()> {
if interface.is_none() {
return Ok(());
}
let loopback = SocketAddr::from((LOCALHOST, coordinator_addr.port()));
if coordinator_addr == loopback {
return Ok(());
}
if std::net::TcpStream::connect_timeout(&loopback, Duration::from_secs(1)).is_ok() {
bail!(
"a coordinator is already running on port {} bound to loopback, so it \
cannot be reached at {coordinator_addr}\n\n \
hint: stop it with `dora down` and re-run, or drop --interface to \
attach to the loopback coordinator",
coordinator_addr.port()
);
}
Ok(())
}
pub(crate) struct CoordinatorStartup {
child: std::process::Child,
stderr: Option<fs::File>,
}
fn start_and_wait_for_coordinator(
coordinator_addr: SocketAddr,
port: u16,
spawn: CoordinatorSpawn,
) -> eyre::Result<crate::ws_client::WsSession> {
let startup = spawn_coordinator(spawn).wrap_err("failed to start dora-coordinator")?;
wait_for_coordinator_start(coordinator_addr, port, startup)
}
pub(crate) fn wait_for_coordinator_start(
coordinator_addr: SocketAddr,
port: u16,
mut startup: CoordinatorStartup,
) -> eyre::Result<crate::ws_client::WsSession> {
let deadline = std::time::Instant::now() + Duration::from_secs(10);
loop {
if let Some(status) = startup
.child
.try_wait()
.wrap_err("failed to poll dora-coordinator")?
{
if let Ok(session) = connect_to_coordinator(coordinator_addr) {
println!("coordinator already running on {coordinator_addr}");
return Ok(session);
}
let stderr = captured_stderr(startup.stderr.as_mut());
let mut message = format!("coordinator exited before it became ready: {status}");
if !stderr.is_empty() {
message.push('\n');
message.push_str(&stderr);
}
if stderr.contains(dora_coordinator::dora_coordinator_store::SCHEMA_MISMATCH_MARKER) {
message.push_str(
"\n\n hint: run `dora up --recreate-store` to archive the incompatible \
store and create a fresh one",
);
}
bail!(message);
}
match connect_to_coordinator(coordinator_addr) {
Ok(session) => return Ok(session),
Err(_) if std::time::Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(50));
}
Err(err) => {
let _ = startup.child.kill();
let _ = startup.child.wait();
bail!(
"timed out waiting for coordinator to start at {coordinator_addr}: {err}\n\n \
hint: is port {port} already in use? Check with `dora status`\n \
or stop the existing coordinator with `dora down`"
);
}
}
}
}
fn captured_stderr(file: Option<&mut fs::File>) -> String {
let Some(file) = file else {
return String::new();
};
let mut bytes = Vec::new();
if file.rewind().is_err() || file.read_to_end(&mut bytes).is_err() {
return String::new();
}
String::from_utf8_lossy(&bytes).trim().to_owned()
}
fn open_coordinator_stderr_capture() -> Option<fs::File> {
let dir = super::coordinator::dora_dir();
fs::create_dir_all(&dir).ok()?;
fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(true)
.open(dir.join("coordinator-stderr.log"))
.ok()
}
#[cfg(not(feature = "redb-backend"))]
const RECREATE_STORE_REQUIRES_REDB: &str = "--recreate-store requires the `redb-backend` feature";
fn archive_default_coordinator_store() -> eyre::Result<()> {
#[cfg(feature = "redb-backend")]
{
let path = super::coordinator::default_redb_path()?;
archive_coordinator_store(&path)
}
#[cfg(not(feature = "redb-backend"))]
{
bail!(RECREATE_STORE_REQUIRES_REDB)
}
}
#[cfg(feature = "redb-backend")]
fn archive_coordinator_store(path: &Path) -> eyre::Result<()> {
use fs2::FileExt as _;
if !path.exists() {
println!(
"no coordinator store at `{}`; nothing to archive",
path.display()
);
return Ok(());
}
let store_file = fs::File::open(path)
.with_context(|| format!("failed to open coordinator store `{}`", path.display()))?;
if store_file.try_lock_exclusive().is_err() {
bail!(
"refusing to archive coordinator store `{}`: another process still \
has it open\n\n \
hint: stop the coordinator with `dora down`, then rerun \
`dora up --recreate-store`",
path.display()
);
}
let backup = available_backup_path(path)?;
fs::rename(path, &backup).with_context(|| {
format!(
"failed to archive coordinator store `{}` as `{}`",
path.display(),
backup.display()
)
})?;
println!(
"archived coordinator store `{}` as `{}`",
path.display(),
backup.display()
);
Ok(())
}
fn lock_default_coordinator_store_recreation() -> eyre::Result<fs::File> {
#[cfg(feature = "redb-backend")]
{
let path = super::coordinator::default_redb_path()?;
lock_coordinator_store_recreation(&path)
}
#[cfg(not(feature = "redb-backend"))]
{
bail!(RECREATE_STORE_REQUIRES_REDB)
}
}
#[cfg(feature = "redb-backend")]
fn lock_coordinator_store_recreation(path: &Path) -> eyre::Result<fs::File> {
use fs2::FileExt as _;
let file_name = path
.file_name()
.context("coordinator store path has no file name")?
.to_string_lossy();
let lock_path = path.with_file_name(format!("{file_name}.recreate.lock"));
let lock = fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&lock_path)
.with_context(|| {
format!(
"failed to open coordinator store recreation lock `{}`",
lock_path.display()
)
})?;
if lock.try_lock_exclusive().is_err() {
println!("waiting for another `dora up --recreate-store` to finish...");
lock.lock_exclusive().with_context(|| {
format!(
"failed to lock coordinator store recreation at `{}`",
lock_path.display()
)
})?;
}
Ok(lock)
}
#[cfg(feature = "redb-backend")]
fn available_backup_path(path: &Path) -> eyre::Result<PathBuf> {
let file_name = path
.file_name()
.context("coordinator store path has no file name")?
.to_string_lossy();
let parent = path
.parent()
.context("coordinator store path has no parent directory")?;
for suffix in 0..u32::MAX {
let backup_name = if suffix == 0 {
format!("{file_name}.backup")
} else {
format!("{file_name}.backup.{suffix}")
};
let candidate = parent.join(backup_name);
if !candidate.exists() {
return Ok(candidate);
}
}
bail!("could not find an available coordinator store backup path")
}
#[derive(Debug)]
pub(crate) enum DownDecision {
Proceed,
Refuse(Vec<DataflowIdAndName>),
}
pub(crate) fn down_decision(running: Vec<DataflowIdAndName>, force: bool) -> DownDecision {
if running.is_empty() || force {
DownDecision::Proceed
} else {
DownDecision::Refuse(running)
}
}
pub(crate) fn down(
config_path: Option<&Path>,
coordinator_addr: SocketAddr,
force: bool,
) -> Result<(), eyre::ErrReport> {
let UpConfig {} = parse_dora_config(config_path)?;
let session = connect_with_retry(coordinator_addr, Duration::from_secs(5)).map_err(|_| {
eyre::eyre!(
"could not connect to coordinator at {coordinator_addr}\n\n \
hint: is it running? Start it with `dora up`"
)
})?;
let running = match send_control_request(&session, &ControlRequest::List) {
Ok(ControlRequestReply::DataflowList(list)) => list.get_active(),
Ok(other) => {
eprintln!(
"warning: could not determine running dataflows before destroying \
(unexpected reply {other:?}); proceeding"
);
Vec::new()
}
Err(err) => {
eprintln!(
"warning: could not determine running dataflows before destroying \
({err}); proceeding"
);
Vec::new()
}
};
if force && !running.is_empty() {
println!(
"Stopping {} running dataflow(s) before teardown...",
running.len()
);
for entry in &running {
if let Err(err) = crate::command::stop::stop_dataflow(entry.uuid, None, true, &session)
{
eprintln!("warning: failed to stop {entry}: {err}");
}
}
}
if let DownDecision::Refuse(running) = down_decision(running, force) {
let list = running
.iter()
.map(|d| format!(" - {d}"))
.collect::<Vec<_>>()
.join("\n");
eyre::bail!(
"refusing to destroy the coordinator at {coordinator_addr}: \
{} dataflow(s) still running\n{list}\n\n \
This tears down the coordinator, its daemons and every dataflow \
above. Lifecycle commands target whatever coordinator owns the \
port, so this may not be the instance you meant — a different \
checkout or another project on this machine shares \
{coordinator_addr} by default.\n\n \
hint: stop the dataflows first (`dora stop --all`), or re-run \
with `--force` if you do mean to kill them\n \
hint: to keep instances isolated, give each one its own port via \
`DORA_COORDINATOR_PORT` (or `--coordinator-port`)",
running.len(),
);
}
println!("Destroying coordinator at {coordinator_addr}");
let reply = send_control_request(&session, &ControlRequest::Destroy)?;
match reply {
ControlRequestReply::DestroyOk => {
println!("Coordinator and daemons destroyed successfully");
}
other => {
bail!("unexpected reply to Destroy: {other:?}");
}
}
Ok(())
}
fn parse_dora_config(config_path: Option<&Path>) -> Result<UpConfig, eyre::ErrReport> {
let path = config_path.or_else(|| Some(Path::new("dora-config.yml")).filter(|p| p.exists()));
let config = match path {
Some(path) => {
let raw = fs::read_to_string(path)
.with_context(|| format!("failed to read `{}`", path.display()))?;
serde_yaml::from_str(&raw)
.with_context(|| format!("failed to parse `{}`", path.display()))?
}
None => Default::default(),
};
Ok(config)
}
static PYTHON_EXECUTABLE_PATH: std::sync::OnceLock<std::ffi::OsString> = std::sync::OnceLock::new();
pub(crate) fn set_python_executable_path(path: std::ffi::OsString) {
let _ = PYTHON_EXECUTABLE_PATH.set(path);
}
fn choose_executable_path(
from_python_wrapper: bool,
recorded: Option<std::ffi::OsString>,
current_exe: impl FnOnce() -> std::io::Result<PathBuf>,
) -> eyre::Result<std::ffi::OsString> {
if from_python_wrapper {
recorded.context(
"could not determine the dora executable path from the Python wrapper \
(sys.argv[0] was not recorded)",
)
} else {
current_exe()
.map(Into::into)
.wrap_err("could not determine dora executable path")
}
}
pub(crate) fn dora_executable_path() -> eyre::Result<std::ffi::OsString> {
choose_executable_path(
cfg!(feature = "python"),
PYTHON_EXECUTABLE_PATH.get().cloned(),
std::env::current_exe,
)
}
#[cfg(test)]
mod executable_path_tests {
use super::choose_executable_path;
use std::ffi::OsString;
use std::path::PathBuf;
#[test]
fn python_wrapper_uses_recorded_argv0() {
let path = choose_executable_path(true, Some(OsString::from("/opt/venv/bin/dora")), || {
panic!("current_exe must not be consulted for the Python wrapper")
})
.expect("recorded path should resolve");
assert_eq!(path, OsString::from("/opt/venv/bin/dora"));
}
#[test]
fn python_wrapper_without_recorded_path_errors() {
let err = choose_executable_path(true, None, || Ok(PathBuf::from("/should/not/be/used")))
.expect_err("missing recorded path must be an error, not a wrong fallback");
assert!(
format!("{err:#}").contains("Python wrapper"),
"unexpected error: {err:#}"
);
}
#[test]
fn standalone_binary_uses_current_exe() {
let path = choose_executable_path(false, None, || Ok(PathBuf::from("/usr/bin/dora")))
.expect("current_exe should resolve");
assert_eq!(path, OsString::from("/usr/bin/dora"));
}
#[test]
fn standalone_binary_ignores_any_recorded_path() {
let path =
choose_executable_path(false, Some(OsString::from("/opt/venv/bin/dora")), || {
Ok(PathBuf::from("/usr/bin/dora"))
})
.expect("current_exe should resolve");
assert_eq!(path, OsString::from("/usr/bin/dora"));
}
}
pub(crate) fn detach_process(cmd: &mut Command) {
detach_process_with_stderr(cmd, Stdio::null());
}
fn detach_process_with_stderr(cmd: &mut Command, stderr: impl Into<Stdio>) {
cmd.stdin(Stdio::null());
cmd.stdout(Stdio::null());
cmd.stderr(stderr);
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
cmd.process_group(0);
}
}
pub(crate) struct CoordinatorSpawn {
pub interface: Option<std::net::IpAddr>,
pub port: Option<u16>,
pub auth: bool,
}
pub(crate) fn spawn_coordinator(spawn: CoordinatorSpawn) -> eyre::Result<CoordinatorStartup> {
let CoordinatorSpawn {
interface,
port,
auth,
} = spawn;
let path = dora_executable_path()?;
let mut cmd = Command::new(path);
cmd.arg("coordinator");
cmd.arg("--quiet");
if let Some(interface) = interface {
cmd.args(["--interface".to_string(), interface.to_string()]);
}
if let Some(port) = port {
cmd.args(["--port".to_string(), port.to_string()]);
}
if auth {
cmd.arg("--auth");
}
let stderr = open_coordinator_stderr_capture();
let child_stderr = stderr
.as_ref()
.and_then(|file| file.try_clone().ok())
.map(Into::into)
.unwrap_or_else(Stdio::null);
detach_process_with_stderr(&mut cmd, child_stderr);
let child = cmd.spawn().wrap_err(
"failed to run `dora coordinator`\n\n \
hint: ensure the `dora` binary is in your PATH",
)?;
match interface {
Some(interface) if !interface.is_loopback() && !interface.is_unspecified() => {
println!("started dora coordinator on {interface} (reachable from other machines)");
println!(
" other commands default to loopback — export \
DORA_COORDINATOR_ADDR={interface} (or pass --coordinator-addr {interface})"
);
}
Some(interface) if !interface.is_loopback() => {
println!("started dora coordinator on {interface} (reachable from other machines)");
}
_ => println!("started dora coordinator"),
}
Ok(CoordinatorStartup { child, stderr })
}
fn start_daemon(coordinator_addr: Option<std::net::IpAddr>, port: u16) -> eyre::Result<()> {
let path = dora_executable_path()?;
let mut cmd = Command::new(path);
cmd.arg("daemon");
cmd.arg("--quiet");
cmd.args(["--coordinator-port".to_string(), port.to_string()]);
if let Some(addr) = coordinator_addr {
cmd.args(["--coordinator-addr".to_string(), addr.to_string()]);
}
detach_process(&mut cmd);
cmd.spawn().wrap_err(
"failed to run `dora daemon`\n\n \
hint: ensure the `dora` binary is in your PATH",
)?;
println!("started dora daemon");
Ok(())
}
#[cfg(test)]
#[cfg(feature = "redb-backend")]
mod tests {
use super::{
LOCALHOST, archive_coordinator_store, available_backup_path, connect_addr_for,
ensure_no_listener_before_archive, lock_coordinator_store_recreation,
};
use std::sync::mpsc;
use std::time::Duration;
#[test]
fn connect_address_follows_a_concrete_bind_interface() {
use std::net::IpAddr;
let iface: IpAddr = "192.168.1.10".parse().unwrap();
assert_eq!(connect_addr_for(None, None).unwrap(), LOCALHOST);
assert_eq!(connect_addr_for(None, Some(iface)).unwrap(), iface);
assert_eq!(
connect_addr_for(None, Some("0.0.0.0".parse().unwrap())).unwrap(),
LOCALHOST
);
let env: IpAddr = "10.9.9.9".parse().unwrap();
assert_eq!(connect_addr_for(Some(env), None).unwrap(), env);
assert_eq!(
connect_addr_for(Some(env), Some("0.0.0.0".parse().unwrap())).unwrap(),
env
);
assert_eq!(connect_addr_for(Some(iface), Some(iface)).unwrap(), iface);
}
#[test]
fn a_bind_interface_contradicting_the_connect_address_is_rejected() {
let err = connect_addr_for(Some(LOCALHOST), Some("192.168.1.10".parse().unwrap()))
.expect_err("contradicting addresses must be refused");
let msg = format!("{err}");
assert!(msg.contains("DORA_COORDINATOR_ADDR"), "{msg}");
assert!(msg.contains("--interface"), "{msg}");
}
#[test]
fn archive_refused_when_a_coordinator_holds_the_port_on_loopback_only() {
let listener =
std::net::TcpListener::bind("127.0.0.1:0").expect("bind loopback coordinator stand-in");
let port = listener.local_addr().expect("read local addr").port();
let unbound: std::net::SocketAddr = format!("10.255.255.1:{port}").parse().unwrap();
let result = ensure_no_listener_before_archive(unbound, eyre::eyre!("connection refused"));
assert!(
result.is_err(),
"archive must be refused while a coordinator holds port {port} on loopback"
);
}
#[test]
fn archive_refused_while_store_file_is_locked() {
use fs2::FileExt as _;
let dir = tempfile::tempdir().expect("create temporary store directory");
let store = dir.path().join("coordinator.redb");
std::fs::write(&store, [1u8]).expect("create store file");
let holder = std::fs::File::open(&store).expect("open store file");
holder
.try_lock_exclusive()
.expect("lock store file like an open database would");
let err = archive_coordinator_store(&store)
.expect_err("archive must be refused while the store file is locked");
assert!(
format!("{err:#}").contains("refusing to archive"),
"unexpected error: {err:#}"
);
assert!(store.exists(), "locked store must not be renamed");
fs2::FileExt::unlock(&holder).expect("release store lock");
archive_coordinator_store(&store).expect("archive must succeed once the lock is free");
assert!(!store.exists(), "store was not renamed");
assert!(
dir.path().join("coordinator.redb.backup").exists(),
"backup was not created"
);
}
#[test]
fn recreate_store_archive_refused_while_port_is_listening() {
let listener = std::net::TcpListener::bind(("127.0.0.1", 0)).expect("bind listener");
let addr = listener.local_addr().expect("listener addr");
let err = ensure_no_listener_before_archive(addr, eyre::eyre!("handshake failed"))
.expect_err("archive must be refused while a listener holds the port");
assert!(
format!("{err:#}").contains("refusing to archive"),
"unexpected error: {err:#}"
);
drop(listener);
ensure_no_listener_before_archive(addr, eyre::eyre!("connection refused"))
.expect("archive must be allowed when nothing is listening");
}
#[test]
fn coordinator_store_backup_path_skips_existing_backups() {
let dir = tempfile::tempdir().expect("create temporary store directory");
let store = dir.path().join("coordinator.redb");
std::fs::write(dir.path().join("coordinator.redb.backup"), []).expect("create backup");
std::fs::write(dir.path().join("coordinator.redb.backup.1"), [])
.expect("create numbered backup");
let backup = available_backup_path(&store).expect("select available backup path");
assert_eq!(backup, dir.path().join("coordinator.redb.backup.2"));
}
#[test]
fn coordinator_store_recreation_lock_serializes_callers() {
let dir = tempfile::tempdir().expect("create temporary store directory");
let store = dir.path().join("coordinator.redb");
let first = lock_coordinator_store_recreation(&store).expect("acquire first lock");
let (tx, rx) = mpsc::channel();
let waiter = std::thread::spawn(move || {
let second = lock_coordinator_store_recreation(&store).expect("acquire second lock");
tx.send(second).expect("report acquired lock");
});
assert!(
rx.recv_timeout(Duration::from_millis(100)).is_err(),
"second recreation acquired the lock before the first released it"
);
drop(first);
let second = rx
.recv_timeout(Duration::from_secs(2))
.expect("second recreation did not acquire the released lock");
drop(second);
waiter.join().expect("join lock waiter");
}
}
#[cfg(test)]
mod coordinator_env_tests {
use super::{DORA_COORDINATOR_PORT_WS_DEFAULT, LOCALHOST, parse_coordinator_env};
use std::net::IpAddr;
#[test]
fn unset_uses_default() {
let addr: IpAddr = parse_coordinator_env("DORA_COORDINATOR_ADDR", None, LOCALHOST).unwrap();
assert_eq!(addr, LOCALHOST);
let port: u16 = parse_coordinator_env(
"DORA_COORDINATOR_PORT",
None,
DORA_COORDINATOR_PORT_WS_DEFAULT,
)
.unwrap();
assert_eq!(port, DORA_COORDINATOR_PORT_WS_DEFAULT);
}
#[test]
fn valid_value_is_parsed() {
let addr: IpAddr =
parse_coordinator_env("DORA_COORDINATOR_ADDR", Some("10.0.0.5"), LOCALHOST).unwrap();
assert_eq!(addr, "10.0.0.5".parse::<IpAddr>().unwrap());
let port: u16 = parse_coordinator_env(
"DORA_COORDINATOR_PORT",
Some("6100"),
DORA_COORDINATOR_PORT_WS_DEFAULT,
)
.unwrap();
assert_eq!(port, 6100);
}
#[test]
fn malformed_value_errors_instead_of_defaulting() {
let err = parse_coordinator_env::<u16>(
"DORA_COORDINATOR_PORT",
Some("not-a-port"),
DORA_COORDINATOR_PORT_WS_DEFAULT,
)
.expect_err("a malformed port must be rejected, not defaulted");
assert!(format!("{err:#}").contains("DORA_COORDINATOR_PORT"));
assert!(
parse_coordinator_env::<IpAddr>("DORA_COORDINATOR_ADDR", Some("999.1.1.1"), LOCALHOST)
.is_err()
);
}
}
#[cfg(test)]
mod destroy_guard_tests {
use super::{DownDecision, down_decision};
use dora_message::coordinator_to_cli::DataflowIdAndName;
fn df(name: &str) -> DataflowIdAndName {
DataflowIdAndName {
uuid: uuid::Uuid::nil(),
name: Some(name.to_string()),
}
}
#[test]
fn down_refuses_when_dataflows_are_running() {
match down_decision(vec![df("long-experiment")], false) {
DownDecision::Refuse(blocked) => {
let names: Vec<_> = blocked.iter().filter_map(|d| d.name.clone()).collect();
assert_eq!(
names,
vec!["long-experiment"],
"the refusal must name what it would have killed, so the \
operator can tell whose coordinator they just hit"
);
}
other => panic!(
"a destroy that would kill running dataflows must stop and say \
so, got {other:?}"
),
}
}
#[test]
fn down_proceeds_when_nothing_is_running() {
assert!(
matches!(down_decision(Vec::new(), false), DownDecision::Proceed),
"an idle coordinator must tear down without a flag"
);
}
#[test]
fn force_destroys_despite_running_dataflows() {
assert!(
matches!(down_decision(vec![df("busy")], true), DownDecision::Proceed),
"`--force` is the documented way to say \"yes, kill them\""
);
}
}