use dora_core::{
descriptor::{CoreNodeKind, CustomNode, Descriptor, DescriptorExt},
topics::{DORA_COORDINATOR_PORT_WS_DEFAULT, LOCALHOST},
types::TypeRegistry,
};
use dora_message::{BuildId, common::GitSource, descriptor::NodeSource, id::NodeId};
use eyre::Context;
use std::{
collections::BTreeMap,
net::IpAddr,
path::{Path, PathBuf},
};
use crate::ws_client::WsSession;
use super::{Executable, default_tracing};
use crate::{
common::{
canonicalize_working_dir, connect_to_coordinator, local_working_dir, resolve_dataflow,
working_dir_or_parent,
},
session::DataflowSession,
};
use distributed::{build_distributed_dataflow, wait_until_dataflow_built};
use local::build_dataflow_locally;
use lockfile::BuildLockfile;
mod distributed;
mod git;
pub mod hub;
mod local;
pub(crate) mod lockfile;
#[derive(Debug, clap::Args)]
pub struct Build {
#[clap(value_name = "PATH")]
dataflow: String,
#[clap(long, value_name = "IP", env = "DORA_COORDINATOR_ADDR")]
coordinator_addr: Option<IpAddr>,
#[clap(long, value_name = "PORT", env = "DORA_COORDINATOR_PORT")]
coordinator_port: Option<u16>,
#[clap(long, action)]
uv: bool,
#[clap(long, action)]
local: bool,
#[clap(long, action)]
strict_types: bool,
#[clap(long, action, conflicts_with = "write_lockfile")]
locked: bool,
#[clap(long, action)]
write_lockfile: bool,
#[clap(long, value_name = "PATH")]
lockfile: Option<PathBuf>,
#[clap(long, action)]
parallel: bool,
#[clap(long, action)]
offline: bool,
#[clap(long = "hub-override", value_name = "PKG=PATH")]
hub_override: Vec<String>,
}
impl Executable for Build {
fn execute(self) -> eyre::Result<()> {
default_tracing()?;
build(BuildConfig {
dataflow: self.dataflow,
coordinator_addr: self.coordinator_addr,
coordinator_port: self.coordinator_port,
uv: self.uv,
force_local: self.local,
strict_types: self.strict_types,
locked: self.locked,
write_lockfile: self.write_lockfile,
lockfile_override: self.lockfile,
parallel: self.parallel,
offline: self.offline,
hub_overrides: self.hub_override,
..Default::default()
})
}
}
#[derive(Debug, Clone, Default)]
pub struct BuildConfig {
pub dataflow: String,
pub coordinator_addr: Option<IpAddr>,
pub coordinator_port: Option<u16>,
pub uv: bool,
pub force_local: bool,
pub strict_types: bool,
pub locked: bool,
pub write_lockfile: bool,
pub lockfile_only: bool,
pub lockfile_override: Option<PathBuf>,
pub parallel: bool,
pub offline: bool,
pub hub_overrides: Vec<String>,
pub working_dir_override: Option<PathBuf>,
}
impl BuildConfig {
pub fn from_str_args(
dataflow: String,
uv: Option<bool>,
coordinator_addr: Option<String>,
coordinator_port: Option<u16>,
force_local: bool,
) -> eyre::Result<Self> {
Ok(Self {
dataflow,
coordinator_addr: coordinator_addr
.map(|addr| addr.parse())
.transpose()
.wrap_err("invalid coordinator_addr")?,
coordinator_port,
uv: uv.unwrap_or_default(),
force_local,
..Default::default()
})
}
}
pub fn build(cfg: BuildConfig) -> eyre::Result<()> {
let BuildConfig {
dataflow,
coordinator_addr,
coordinator_port,
uv,
force_local,
strict_types,
locked,
write_lockfile,
lockfile_only,
lockfile_override,
parallel,
offline,
hub_overrides,
working_dir_override,
} = cfg;
let mut hub_override_dirs: BTreeMap<String, PathBuf> = BTreeMap::new();
for spec in &hub_overrides {
let (pkg, path) = spec.split_once('=').ok_or_else(|| {
eyre::eyre!("invalid --hub-override `{spec}`: expected `<namespace>/<name>=<path>`")
})?;
let key = dora_hub_client::reference::PackageRef::parse(pkg.trim())
.with_context(|| format!("invalid --hub-override package `{pkg}`"))?
.key();
let dir = std::fs::canonicalize(path.trim()).with_context(|| {
format!("invalid --hub-override path `{}` for `{pkg}`", path.trim())
})?;
hub_override_dirs.insert(key, dir);
}
if dataflow.is_empty() {
eyre::bail!(
"BuildConfig::dataflow is empty — set it to a YAML path or URL before calling build()"
);
}
let dataflow_path = resolve_dataflow(dataflow).context("could not resolve dataflow")?;
if lockfile_override.is_some() && !(locked || write_lockfile) {
eyre::bail!("`--lockfile` requires either `--locked` or `--write-lockfile`");
}
let working_dir = working_dir_or_parent(working_dir_override.as_deref(), &dataflow_path);
let mut dataflow_descriptor = Descriptor::blocking_read(&dataflow_path)
.wrap_err_with(|| {
format!(
"failed to read dataflow at `{}`\n\n \
hint: check the file exists and is valid YAML",
dataflow_path.display()
)
})?
.expand(working_dir)
.wrap_err("failed to expand modules in dataflow descriptor")?;
if !hub_override_dirs.is_empty() {
if coordinator_addr.is_some() || coordinator_port.is_some() {
eyre::bail!(
"`--hub-override` is a local build feature and cannot be combined with a remote \
coordinator (`--coordinator-addr`/`--coordinator-port`)"
);
}
if !force_local && dataflow_descriptor.nodes.iter().any(|n| n.deploy.is_some()) {
eyre::bail!(
"`--hub-override` is a local build feature and cannot be used with a distributed \
(`deploy:`) dataflow — the local checkout only exists on this machine. Use \
`--local` to force a fully local build if that is what you want."
);
}
if write_lockfile {
eyre::bail!(
"`--hub-override` cannot be combined with `--write-lockfile`: the override \
substitutes local source for a hub node, so the regenerated lockfile would drop \
that node's hub pin. Drop `--write-lockfile` (or the override) when refreshing \
the lockfile."
);
}
}
let has_hub_nodes = dataflow_descriptor.nodes.iter().any(|n| n.hub.is_some());
let source_fingerprint = has_hub_nodes
.then(|| DataflowSession::fingerprint_source(&dataflow_descriptor))
.flatten();
let strict = strict_types || dataflow_descriptor.strict_types.unwrap_or(false);
let mut registry = TypeRegistry::new();
let types_dir = working_dir.join("types");
if types_dir.is_dir() {
match registry.load_from_dir(&types_dir) {
Ok(count) if count > 0 => {
log::info!("Loaded {count} user-defined type(s) from types/");
}
Err(e) => {
eyre::bail!("failed to load user types: {e}");
}
_ => {}
}
}
let lockfile_path = BuildLockfile::path_for_dataflow(&dataflow_path, lockfile_override);
let build_lockfile = if locked {
Some(BuildLockfile::read_from(&lockfile_path).with_context(|| {
format!(
"failed to read build lockfile at `{}`",
lockfile_path.display()
)
})?)
} else {
None
};
let hub_pins = build_lockfile.as_ref().map(|l| l.git_sources.clone());
let hub_binary_pins = build_lockfile.as_ref().map(|l| l.binary_sources.clone());
let hub_resolution = hub::resolve_hub_nodes(
&mut dataflow_descriptor,
&mut registry,
offline,
hub_pins.as_ref(),
hub_binary_pins.as_ref(),
locked,
&hub_override_dirs,
)?;
let hub_override_node_dirs = hub_resolution.override_dirs.clone();
for note in &hub_resolution.notes {
println!(" {note}");
}
for warning in &hub_resolution.warnings {
eprintln!(" warning: {warning}");
}
let resolved_dataflow_for_session =
(!hub_resolution.is_empty()).then(|| dataflow_descriptor.clone());
let injection = dora_core::manifest::inject::inject_adjacent_manifests(
&mut dataflow_descriptor,
working_dir,
&mut registry,
);
for note in &injection.notes {
println!(" {note}");
}
let type_result = dora_core::descriptor::validate::check_type_annotations_full(
&dataflow_descriptor,
®istry,
strict,
);
for inf in &type_result.inferences {
println!(" {inf}");
}
let warning_count = injection.warnings.len() + type_result.warnings.len();
if warning_count > 0 {
for w in &injection.warnings {
eprintln!(" warning: {w}");
}
for w in &type_result.warnings {
eprintln!(" warning: {w}");
}
if strict {
eyre::bail!("{warning_count} type error(s) found (strict mode)");
} else {
eprintln!(
"{warning_count} type warning(s) found.\n \
hint: use --strict-types to fail on type warnings"
);
}
}
let mut git_sources = BTreeMap::new();
let mut descriptor_git_sources = BTreeMap::new();
let resolved_nodes = dataflow_descriptor
.resolve_aliases_and_set_defaults()
.context("failed to resolve nodes")?;
let session_build_fingerprint = DataflowSession::fingerprint_build_inputs(&resolved_nodes);
for (node_id, node) in &resolved_nodes {
if let CoreNodeKind::Custom(CustomNode {
source: NodeSource::GitBranch { repo, rev },
..
}) = &node.kind
{
descriptor_git_sources.insert(
node_id.clone(),
NodeSource::GitBranch {
repo: repo.clone(),
rev: rev.clone(),
},
);
}
}
let descriptor_fingerprint =
BuildLockfile::fingerprint_descriptor_git_sources(&descriptor_git_sources);
if let Some(lockfile) = &build_lockfile {
lockfile
.ensure_descriptor_fingerprint_matches(&descriptor_fingerprint)
.with_context(|| {
format!(
"failed to validate lockfile against descriptor at `{}`",
dataflow_path.display()
)
})?;
}
for (node_id, node) in resolved_nodes {
if let CoreNodeKind::Custom(CustomNode {
source: NodeSource::GitBranch { repo, rev },
..
}) = node.kind
{
let source = match hub_resolution.sources.get(&node_id) {
Some(source) => source.clone(),
None => match &build_lockfile {
Some(lockfile) => lockfile.get_source(&node_id, &repo).with_context(|| {
format!("failed to resolve locked git source `{node_id}`")
})?,
None => git::fetch_commit_hash(repo, rev)
.with_context(|| format!("failed to find commit hash for `{node_id}`"))?,
},
};
git_sources.insert(node_id, source);
}
}
if write_lockfile {
BuildLockfile::write_git_sources(
&lockfile_path,
&git_sources,
&hub_resolution.binary_sources,
&descriptor_fingerprint,
)
.with_context(|| {
format!(
"failed to write build lockfile to `{}`",
lockfile_path.display()
)
})?;
log::info!("wrote build lockfile to {}", lockfile_path.display());
}
if lockfile_only {
return Ok(());
}
let mut dataflow_session =
DataflowSession::read_session(&dataflow_path).context("failed to read DataflowSession")?;
let session = || connect_to_coordinator_with_defaults(coordinator_addr, coordinator_port);
let build_kind = if !hub_override_dirs.is_empty() {
log::info!("Building locally because `--hub-override` was given");
BuildKind::Local
} else if force_local {
log::info!("Building locally, as requested through `--force-local`");
BuildKind::Local
} else if dataflow_descriptor.nodes.iter().all(|n| n.deploy.is_none()) {
log::info!("Building locally because dataflow does not contain any `deploy` sections");
BuildKind::Local
} else if coordinator_addr.is_some() || coordinator_port.is_some() {
log::info!("Building through coordinator, using the given coordinator socket information");
BuildKind::ThroughCoordinator {
coordinator_session: session().context("failed to connect to coordinator")?,
}
} else {
match session() {
Ok(coordinator_session) => {
log::info!("Found local dora coordinator instance -> building through coordinator");
BuildKind::ThroughCoordinator {
coordinator_session,
}
}
Err(_) => {
log::warn!("No dora coordinator instance found -> trying a local build");
BuildKind::Local
}
}
};
match build_kind {
BuildKind::Local => {
log::info!("running local build");
let local_working_dir =
canonicalize_working_dir(working_dir_override.as_deref(), &dataflow_path)?;
let build_info = build_dataflow_locally(
dataflow_descriptor,
&git_sources,
&dataflow_session,
local_working_dir,
uv,
parallel,
&hub_override_node_dirs,
)?;
dataflow_session.git_sources = git_sources;
if dataflow_session.build_id.is_none() {
dataflow_session.build_id = Some(BuildId::generate());
}
dataflow_session.local_build = Some(build_info);
dataflow_session.build_fingerprint = Some(session_build_fingerprint.clone());
dataflow_session.resolved_dataflow = resolved_dataflow_for_session.clone();
dataflow_session.source_fingerprint = source_fingerprint.clone();
dataflow_session
.write_out_for_dataflow(&dataflow_path)
.context("failed to write out dataflow session file")?;
}
BuildKind::ThroughCoordinator {
coordinator_session,
} => {
let inferred_local_working_dir =
local_working_dir(&dataflow_path, &dataflow_descriptor, &coordinator_session)?;
let local_working_dir = select_distributed_working_dir(
working_dir_override.as_deref(),
inferred_local_working_dir,
&dataflow_path,
)?;
let build_id = build_distributed_dataflow(
&coordinator_session,
dataflow_descriptor,
&git_sources,
&dataflow_session,
local_working_dir,
uv,
)?;
let build_result =
wait_until_dataflow_built(build_id, &coordinator_session, log::LevelFilter::Info);
dataflow_session.resolved_dataflow = resolved_dataflow_for_session.clone();
dataflow_session.source_fingerprint = source_fingerprint.clone();
finalize_distributed_build_session(
&mut dataflow_session,
&dataflow_path,
git_sources,
build_result,
session_build_fingerprint,
)?;
}
};
Ok(())
}
enum BuildKind {
Local,
ThroughCoordinator { coordinator_session: WsSession },
}
fn connect_to_coordinator_with_defaults(
coordinator_addr: Option<std::net::IpAddr>,
coordinator_port: Option<u16>,
) -> eyre::Result<WsSession> {
let coordinator_addr = coordinator_addr.unwrap_or(LOCALHOST);
let coordinator_port = coordinator_port.unwrap_or(DORA_COORDINATOR_PORT_WS_DEFAULT);
connect_to_coordinator((coordinator_addr, coordinator_port).into())
}
fn select_distributed_working_dir(
working_dir_override: Option<&Path>,
inferred_local_working_dir: Option<PathBuf>,
dataflow_path: &Path,
) -> eyre::Result<Option<PathBuf>> {
match (working_dir_override, inferred_local_working_dir) {
(Some(override_), Some(_)) => {
let canonical = canonicalize_working_dir(Some(override_), dataflow_path)?;
Ok(Some(canonical))
}
(Some(_), None) => eyre::bail!(
"`working_dir_override` can only be used for single-machine coordinator builds where CLI and daemon run on the same machine"
),
(None, inferred) => Ok(inferred),
}
}
fn finalize_distributed_build_session(
dataflow_session: &mut DataflowSession,
dataflow_path: &Path,
git_sources: BTreeMap<NodeId, GitSource>,
build_result: eyre::Result<BuildId>,
session_build_fingerprint: String,
) -> eyre::Result<()> {
let build_id = build_result?;
dataflow_session.git_sources = git_sources;
dataflow_session.build_id = Some(build_id);
dataflow_session.local_build = None;
dataflow_session.build_fingerprint = Some(session_build_fingerprint);
dataflow_session
.write_out_for_dataflow(dataflow_path)
.context("failed to write out dataflow session file")?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn from_str_args_parses_valid_coordinator_addr() {
let cfg = BuildConfig::from_str_args(
"dataflow.yml".into(),
None,
Some("127.0.0.1".into()),
None,
false,
)
.expect("valid IP should parse");
assert_eq!(
cfg.coordinator_addr,
Some("127.0.0.1".parse::<IpAddr>().unwrap())
);
}
#[test]
fn from_str_args_accepts_none_addr() {
let cfg = BuildConfig::from_str_args("dataflow.yml".into(), None, None, None, false)
.expect("None addr should be fine");
assert!(cfg.coordinator_addr.is_none());
}
#[test]
fn from_str_args_errors_on_invalid_addr() {
let err = BuildConfig::from_str_args(
"dataflow.yml".into(),
None,
Some("not-an-ip".into()),
None,
false,
)
.expect_err("malformed addr should error");
assert!(
err.to_string().contains("invalid coordinator_addr"),
"error should carry context: {err}"
);
}
#[test]
fn from_str_args_unwraps_uv_default_to_false() {
let cfg =
BuildConfig::from_str_args("dataflow.yml".into(), None, None, None, false).unwrap();
assert!(!cfg.uv, "uv should default to false when None");
}
fn unique_temp_path(name: &str) -> PathBuf {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
std::env::temp_dir().join(format!(
"dora-build-mod-tests-{name}-{}-{nanos}",
std::process::id()
))
}
#[test]
fn distributed_override_is_used_when_local_working_dir_allowed() {
let root = unique_temp_path("override-ok");
std::fs::create_dir_all(&root).unwrap();
let selected = select_distributed_working_dir(
Some(root.as_path()),
Some(PathBuf::from("/tmp/inferred")),
Path::new("/tmp/dataflow.yml"),
)
.unwrap();
assert_eq!(selected, Some(dunce::canonicalize(&root).unwrap()));
std::fs::remove_dir_all(root).unwrap();
}
#[test]
fn distributed_override_is_rejected_for_non_local_builds() {
let root = unique_temp_path("override-reject");
std::fs::create_dir_all(&root).unwrap();
let err = select_distributed_working_dir(
Some(root.as_path()),
None,
Path::new("/tmp/dataflow.yml"),
)
.unwrap_err();
assert!(
err.to_string()
.contains("can only be used for single-machine coordinator builds")
);
std::fs::remove_dir_all(root).unwrap();
}
#[test]
fn distributed_without_override_uses_inferred_working_dir() {
let inferred = Some(PathBuf::from("/tmp/inferred"));
let selected =
select_distributed_working_dir(None, inferred.clone(), Path::new("/tmp/dataflow.yml"))
.unwrap();
assert_eq!(selected, inferred);
}
fn git_source(repo: &str, commit_hash: &str) -> GitSource {
GitSource {
repo: repo.to_owned(),
commit_hash: commit_hash.to_owned(),
subdir: None,
hub: None,
}
}
#[test]
fn distributed_build_failure_does_not_persist_new_git_sources() {
let tmp = tempfile::tempdir().unwrap();
let dataflow_path = tmp.path().join("dataflow.yml");
std::fs::write(&dataflow_path, "nodes: []\n").unwrap();
let old_build_id = BuildId::generate();
let old_git_sources = BTreeMap::from([(
"old-node".parse().unwrap(),
git_source("https://example.com/old.git", "old-commit"),
)]);
let initial_session = DataflowSession {
build_id: Some(old_build_id),
git_sources: old_git_sources.clone(),
build_fingerprint: Some("old-fingerprint".to_owned()),
..DataflowSession::default()
};
let mut in_memory_session = initial_session.clone();
initial_session
.write_out_for_dataflow(&dataflow_path)
.unwrap();
let new_git_sources = BTreeMap::from([(
"new-node".parse().unwrap(),
git_source("https://example.com/new.git", "new-commit"),
)]);
let err = finalize_distributed_build_session(
&mut in_memory_session,
&dataflow_path,
new_git_sources,
Err(eyre::eyre!("coordinator build failed")),
"new-fingerprint".to_owned(),
)
.unwrap_err();
assert!(err.to_string().contains("coordinator build failed"));
assert_eq!(in_memory_session.build_id, Some(old_build_id));
assert_eq!(in_memory_session.git_sources, old_git_sources);
assert_eq!(
in_memory_session.build_fingerprint.as_deref(),
Some("old-fingerprint")
);
let persisted = DataflowSession::read_session(&dataflow_path).unwrap();
assert_eq!(persisted.build_id, initial_session.build_id);
assert_eq!(persisted.git_sources, initial_session.git_sources);
assert_eq!(
persisted.build_fingerprint,
initial_session.build_fingerprint
);
}
#[test]
fn distributed_build_success_persists_git_sources_and_build_id() {
let tmp = tempfile::tempdir().unwrap();
let dataflow_path = tmp.path().join("dataflow.yml");
std::fs::write(&dataflow_path, "nodes: []\n").unwrap();
let mut session = DataflowSession::default();
session.write_out_for_dataflow(&dataflow_path).unwrap();
let build_id = BuildId::generate();
let git_sources = BTreeMap::from([(
"node".parse().unwrap(),
git_source("https://example.com/repo.git", "abc123"),
)]);
finalize_distributed_build_session(
&mut session,
&dataflow_path,
git_sources.clone(),
Ok(build_id),
"new-fingerprint".to_owned(),
)
.unwrap();
assert_eq!(session.build_id, Some(build_id));
assert_eq!(session.git_sources, git_sources);
assert!(session.local_build.is_none());
assert_eq!(
session.build_fingerprint.as_deref(),
Some("new-fingerprint")
);
let persisted = DataflowSession::read_session(&dataflow_path).unwrap();
assert_eq!(persisted.build_id, Some(build_id));
assert_eq!(persisted.git_sources, session.git_sources);
assert!(persisted.local_build.is_none());
assert_eq!(
persisted.build_fingerprint.as_deref(),
Some("new-fingerprint")
);
}
}