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)]
recreate_store: bool,
}
impl Executable for Up {
fn execute(self) -> eyre::Result<()> {
default_tracing()?;
up(self.config.as_deref(), self.auth, self.recreate_store)
}
}
#[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) -> eyre::Result<()> {
let UpConfig {} = parse_dora_config(config_path)?;
let addr: std::net::IpAddr = coordinator_env_value("DORA_COORDINATOR_ADDR", LOCALHOST)?;
let port: u16 =
coordinator_env_value("DORA_COORDINATOR_PORT", DORA_COORDINATOR_PORT_WS_DEFAULT)?;
let coordinator_addr = (addr, port).into();
if recreate_store && !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 {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, auth)?
}
}
}
Err(_) => start_and_wait_for_coordinator(coordinator_addr, port, auth)?,
};
if !daemon_running(&session)? {
start_daemon().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 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,
{
match raw {
Some(s) => s
.parse()
.map_err(|e| eyre::eyre!("invalid {var_name}: {s:?} ({e})")),
None => Ok(default),
}
}
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<()> {
if std::net::TcpStream::connect_timeout(&coordinator_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 \
{coordinator_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(())
}
struct CoordinatorStartup {
child: std::process::Child,
stderr: Option<fs::File>,
}
fn start_and_wait_for_coordinator(
coordinator_addr: SocketAddr,
port: u16,
auth: bool,
) -> eyre::Result<crate::ws_client::WsSession> {
let startup = start_coordinator(auth).wrap_err("failed to start dora-coordinator")?;
wait_for_coordinator_start(coordinator_addr, port, startup)
}
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)
}
pub(crate) fn dora_executable_path() -> eyre::Result<std::ffi::OsString> {
if cfg!(feature = "python") {
std::env::args_os()
.nth(1)
.context("could not get dora path from Python wrapper arguments")
} else {
std::env::current_exe()
.map(Into::into)
.wrap_err("could not determine dora executable path")
}
}
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);
}
}
fn start_coordinator(auth: bool) -> eyre::Result<CoordinatorStartup> {
let path = dora_executable_path()?;
let mut cmd = Command::new(path);
cmd.arg("coordinator");
cmd.arg("--quiet");
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",
)?;
println!("started dora coordinator");
Ok(CoordinatorStartup { child, stderr })
}
fn start_daemon() -> eyre::Result<()> {
let path = dora_executable_path()?;
let mut cmd = Command::new(path);
cmd.arg("daemon");
cmd.arg("--quiet");
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::{
archive_coordinator_store, available_backup_path, ensure_no_listener_before_archive,
lock_coordinator_store_recreation,
};
use std::sync::mpsc;
use std::time::Duration;
#[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\""
);
}
}