use super::{Executable, default_tracing};
use crate::{
command::start::attach::attach_dataflow,
common::{
CoordinatorOptions, connect_to_coordinator, error_indicates_dataflow_finished,
expect_reply, local_working_dir, resolve_dataflow, send_control_request, write_events_to,
},
output::{LogOutputConfig, print_log_message},
session::DataflowSession,
ws_client::WsSession,
};
use dora_core::descriptor::{Descriptor, DescriptorExt};
use dora_message::{cli_to_coordinator::ControlRequest, common::LogMessage, descriptor::EnvValue};
use eyre::Context;
use std::{
collections::BTreeMap,
io::IsTerminal,
net::SocketAddr,
path::{Path, PathBuf},
};
use uuid::Uuid;
mod attach;
#[derive(Debug, clap::Args)]
pub struct Start {
#[clap(value_name = "PATH")]
dataflow: String,
#[clap(long, short = 'n')]
name: Option<String>,
#[clap(flatten)]
coordinator: CoordinatorOptions,
#[clap(long, action, conflicts_with = "detach")]
attach: bool,
#[clap(long, action)]
detach: bool,
#[clap(long, action)]
hot_reload: bool,
#[clap(long, action)]
uv: bool,
#[clap(long, action)]
debug: bool,
#[clap(
long,
num_args = 0..=1,
require_equals = true,
default_missing_value = "true",
value_name = "BOOL"
)]
pub exit_when_nodes_finish: Option<bool>,
#[clap(long = "env", value_name = "KEY=VALUE")]
env: Vec<String>,
}
impl Executable for Start {
fn execute(self) -> eyre::Result<()> {
default_tracing()?;
let coordinator_socket = self.coordinator.socket_addr();
let env_overrides = crate::env_overrides::parse_env_overrides(&self.env)?;
let (dataflow, dataflow_descriptor, session, dataflow_id) = start_dataflow(
self.dataflow,
self.name,
coordinator_socket,
self.uv,
self.debug,
env_overrides,
self.exit_when_nodes_finish,
)?;
let attach = match (self.attach, self.detach) {
(true, _) => true,
(false, true) => false,
(false, false) => {
if std::io::stdin().is_terminal() {
eprintln!("attaching to dataflow (use `--detach` to run in background)");
true
} else {
eprintln!("non-interactive mode: running in background");
false
}
}
};
if attach {
let log_level = env_logger::Builder::new()
.filter_level(log::LevelFilter::Info)
.parse_default_env()
.build()
.filter();
attach_dataflow(
dataflow_descriptor,
dataflow,
dataflow_id,
&session,
self.hot_reload,
log_level,
)
} else {
let print_daemon_name = dataflow_descriptor.nodes.iter().any(|n| n.deploy.is_some());
wait_until_dataflow_started(
dataflow_id,
&session,
log::LevelFilter::Info,
print_daemon_name,
)
}
}
}
fn start_dataflow(
dataflow: String,
name: Option<String>,
coordinator_socket: SocketAddr,
uv: bool,
debug: bool,
env_overrides: BTreeMap<String, EnvValue>,
exit_when_nodes_finish: Option<bool>,
) -> Result<(PathBuf, Descriptor, WsSession, Uuid), eyre::Error> {
let dataflow = resolve_dataflow(dataflow).context("could not resolve dataflow")?;
let (dataflow_descriptor, dataflow_session) =
prepare_descriptor(&dataflow, debug, env_overrides, exit_when_nodes_finish)?;
let session = connect_to_coordinator(coordinator_socket)?;
let local_working_dir = local_working_dir(&dataflow, &dataflow_descriptor, &session)?;
let dataflow_id = {
let dataflow = dataflow_descriptor.clone();
let reply = send_control_request(
&session,
&ControlRequest::Start {
build_id: dataflow_session.build_id,
session_id: dataflow_session.session_id,
dataflow,
name,
local_working_dir,
uv,
write_events_to: write_events_to(),
},
)?;
let uuid = expect_reply!(reply, DataflowStartTriggered { uuid })?;
println!("dataflow start triggered: {uuid}");
uuid
};
Ok((dataflow, dataflow_descriptor, session, dataflow_id))
}
fn prepare_descriptor(
dataflow: &Path,
debug: bool,
env_overrides: BTreeMap<String, EnvValue>,
exit_when_nodes_finish: Option<bool>,
) -> eyre::Result<(Descriptor, DataflowSession)> {
let working_dir = dataflow
.parent()
.filter(|p| !p.as_os_str().is_empty())
.unwrap_or_else(|| std::path::Path::new("."));
let mut dataflow_descriptor = Descriptor::blocking_read(dataflow)
.wrap_err_with(|| {
format!(
"failed to read dataflow at `{}`\n\n \
hint: check the file exists, is valid YAML, and matches the dataflow schema (see details below)",
dataflow.display()
)
})?
.expand(working_dir)
.wrap_err("failed to expand modules in dataflow descriptor")?;
let mut dataflow_session =
DataflowSession::read_session(dataflow).context("failed to read DataflowSession")?;
if dataflow_descriptor.nodes.iter().any(|n| n.hub.is_some()) {
let resolved = dataflow_session.resolved_dataflow.clone().ok_or_else(|| {
eyre::eyre!("this dataflow uses `hub:` nodes — run `dora build` first")
})?;
let current = DataflowSession::fingerprint_source(&dataflow_descriptor);
if current.is_none() || current != dataflow_session.source_fingerprint {
eyre::bail!(
"this dataflow changed since the last `dora build` — run `dora build` again \
(`dora start` cannot re-resolve `hub:` references)"
);
}
dataflow_descriptor = resolved;
}
let resolved_for_fingerprint = dataflow_descriptor
.resolve_aliases_and_set_defaults()
.context("failed to resolve nodes for session fingerprint")?;
if dataflow_session.invalidate_if_build_inputs_changed(&resolved_for_fingerprint) {
dataflow_session
.write_out_for_dataflow(dataflow)
.context("failed to persist invalidated dataflow session")?;
}
drop(resolved_for_fingerprint);
let has_deploy_nodes = dataflow_descriptor.nodes.iter().any(|n| n.deploy.is_some());
if local_build_blocks_distributed_start(
dataflow_session.local_build.is_some(),
has_deploy_nodes,
) {
eyre::bail!(
"this dataflow was built locally, but it has `deploy` sections and is started \
through the coordinator — remote daemons cannot use a local build.\n\n \
run `dora build` against the running coordinator (without `--local`) to build on \
the target machines before `dora start`"
);
}
if debug {
dataflow_descriptor.debug.enable_debug_inspection = true;
}
crate::env_overrides::apply_env_overrides(&mut dataflow_descriptor, env_overrides);
dataflow_descriptor.apply_exit_when_nodes_finish(exit_when_nodes_finish);
Ok((dataflow_descriptor, dataflow_session))
}
fn wait_until_dataflow_started(
dataflow_id: Uuid,
session: &WsSession,
log_level: log::LevelFilter,
print_daemon_id: bool,
) -> eyre::Result<()> {
match session.subscribe_logs(
&serde_json::to_vec(&ControlRequest::LogSubscribe {
dataflow_id,
level: log_level,
})
.wrap_err("failed to serialize message")?,
) {
Ok(log_rx) => {
std::thread::spawn(move || {
while let Ok(Ok(raw)) = log_rx.recv() {
let parsed: eyre::Result<LogMessage> =
serde_json::from_slice(&raw).context("failed to parse log message");
match parsed {
Ok(log_message) => {
let config = LogOutputConfig {
print_daemon_name: print_daemon_id,
..LogOutputConfig::default()
};
print_log_message(log_message, &config);
}
Err(err) => {
tracing::warn!("failed to parse log message: {err:?}")
}
}
}
});
}
Err(err) if error_indicates_dataflow_finished(&err.to_string()) => {
tracing::debug!("dataflow {dataflow_id} completed before log subscribe arrived");
}
Err(err) => return Err(err).wrap_err("failed to subscribe to logs"),
}
match send_control_request(session, &ControlRequest::WaitForSpawn { dataflow_id }) {
Ok(reply) => {
let uuid = expect_reply!(reply, DataflowSpawned { uuid })?;
println!("dataflow started: {uuid}");
}
Err(err) if error_indicates_dataflow_finished(&err.to_string()) => {
println!("dataflow started and finished: {dataflow_id}");
}
Err(err) => {
return Err(err).wrap_err(
"dataflow failed to start\n\n \
hint: if nodes require building, run `dora build <dataflow.yml>` first",
);
}
}
Ok(())
}
fn local_build_blocks_distributed_start(has_local_build: bool, has_deploy_nodes: bool) -> bool {
has_local_build && has_deploy_nodes
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn local_build_blocks_distributed_start_only_when_deployed_and_local() {
assert!(local_build_blocks_distributed_start(true, true));
assert!(!local_build_blocks_distributed_start(false, true));
assert!(!local_build_blocks_distributed_start(true, false));
assert!(!local_build_blocks_distributed_start(false, false));
}
fn hub_dataflow_fixture(dir: &Path, source: &str, resolved: &str) -> PathBuf {
std::fs::create_dir_all(dir).unwrap();
let dataflow = dir.join("hubflow.yml");
std::fs::write(&dataflow, source).unwrap();
let expanded = Descriptor::parse(source.as_bytes().to_vec())
.unwrap()
.expand(dir)
.unwrap();
let resolved = Descriptor::parse(resolved.as_bytes().to_vec()).unwrap();
let session = DataflowSession {
source_fingerprint: DataflowSession::fingerprint_source(&expanded),
resolved_dataflow: Some(resolved.clone()),
build_fingerprint: Some(DataflowSession::fingerprint_build_inputs(
&resolved.resolve_aliases_and_set_defaults().unwrap(),
)),
..Default::default()
};
session.write_out_for_dataflow(&dataflow).unwrap();
dataflow
}
const HUB_SOURCE: &str = "nodes:\n - id: a\n hub: dora-yolo@^0.5\n";
const HUB_RESOLVED: &str = "nodes:\n - id: a\n path: ./a\n";
#[test]
fn hub_start_applies_exit_when_nodes_finish_in_both_directions() {
let dir = tempfile::tempdir().unwrap();
let dataflow = hub_dataflow_fixture(dir.path(), HUB_SOURCE, HUB_RESOLVED);
let (descriptor, _) =
prepare_descriptor(&dataflow, false, BTreeMap::new(), Some(true)).unwrap();
assert_eq!(descriptor.exit_when_nodes_finish, Some(true));
assert!(descriptor.nodes.iter().all(|n| n.hub.is_none()));
let (descriptor, _) = prepare_descriptor(&dataflow, false, BTreeMap::new(), None).unwrap();
assert_eq!(descriptor.exit_when_nodes_finish, None);
let on = "exit_when_nodes_finish: true\n";
let dataflow = hub_dataflow_fixture(
&dir.path().join("forced"),
&format!("{on}{HUB_SOURCE}"),
&format!("{on}{HUB_RESOLVED}"),
);
let (descriptor, _) =
prepare_descriptor(&dataflow, false, BTreeMap::new(), Some(false)).unwrap();
assert_eq!(descriptor.exit_when_nodes_finish, Some(false));
let (descriptor, _) = prepare_descriptor(&dataflow, false, BTreeMap::new(), None).unwrap();
assert_eq!(descriptor.exit_when_nodes_finish, Some(true));
}
#[test]
fn hub_start_still_rejects_a_dataflow_edited_since_the_build() {
let dir = tempfile::tempdir().unwrap();
let dataflow = hub_dataflow_fixture(dir.path(), HUB_SOURCE, HUB_RESOLVED);
std::fs::write(&dataflow, format!("{HUB_SOURCE} - id: b\n path: ./b\n")).unwrap();
let err = prepare_descriptor(&dataflow, false, BTreeMap::new(), Some(true)).unwrap_err();
assert!(
err.to_string()
.contains("changed since the last `dora build`"),
"expected the staleness bail, got: {err:#}"
);
}
}