Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions binaries/cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ redb-backend = ["dora-coordinator/redb-backend", "dep:fs2", "dep:libc"]
[dependencies]
arrow = { workspace = true, features = ["ipc"] }
arrow-schema = { workspace = true }
base64 = "0.22"
clap = { version = "4.6.1", features = ["derive", "string", "env"] }
clap_complete = "4.6.7"
eyre = { workspace = true }
Expand Down
34 changes: 32 additions & 2 deletions binaries/cli/src/command/record.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ use std::{
time::SystemTime,
};

use base64::prelude::{BASE64_STANDARD, Engine as _};
use clap::Args;
use dora_core::descriptor::Descriptor;
use dora_message::{
Expand Down Expand Up @@ -292,6 +293,16 @@ fn build_record_inputs(topic_map: &BTreeMap<&str, &str>, queue_size: u64) -> ser
inputs
}

/// The record node's env var carrying the original descriptor, base64-encoded.
///
/// Every `env:` value goes through `$VAR` expansion each time the
/// descriptor is parsed (by the CLI, then again by the daemon), so the raw
/// YAML must not travel as one: a `$NAME` anywhere in it, even in a
/// comment, failed `dora record` when unset, and was substituted otherwise,
/// writing e.g. the value of an `${API_TOKEN}` reference into the
/// recording's header. The base64 alphabet has no `$`.
const DESCRIPTOR_ENV: &str = "DORA_RECORD_DESCRIPTOR_BASE64";

fn run_record(args: Record) -> eyre::Result<()> {
let yaml_bytes =
std::fs::read(&args.file).wrap_err_with(|| format!("failed to read {}", args.file))?;
Expand Down Expand Up @@ -364,8 +375,8 @@ fn run_record(args: Record) -> eyre::Result<()> {
serde_yaml::Value::String(topics_json),
);
env_mapping.insert(
serde_yaml::Value::String("DORA_RECORD_DESCRIPTOR".to_string()),
serde_yaml::Value::String(String::from_utf8_lossy(&yaml_bytes).to_string()),
serde_yaml::Value::String(DESCRIPTOR_ENV.to_string()),
serde_yaml::Value::String(BASE64_STANDARD.encode(&yaml_bytes)),
);

// Build the record node YAML entry
Expand Down Expand Up @@ -823,6 +834,25 @@ mod tests {
.collect()
}

/// The descriptor reaches the record node through an `env:` value, which
/// is `$VAR`-expanded on every parse. Its encoding must survive that
/// unchanged, whatever `$` references the YAML contains.
#[test]
fn descriptor_env_value_survives_env_expansion() {
let yaml = b"# set $DORA_TEST_SURELY_UNSET_VAR first\nnodes:\n- id: a\n env:\n TOKEN: ${HOME}\n";
let env = serde_yaml::Mapping::from_iter([(
serde_yaml::Value::String(DESCRIPTOR_ENV.to_string()),
serde_yaml::Value::String(BASE64_STANDARD.encode(yaml)),
)]);
let parsed: BTreeMap<String, dora_message::descriptor::EnvValue> =
serde_yaml::from_value(serde_yaml::Value::Mapping(env))
.expect("the encoded descriptor must parse as an env value");
let decoded = BASE64_STANDARD
.decode(parsed[DESCRIPTOR_ENV].to_string())
.expect("valid base64");
assert_eq!(decoded, yaml);
}

#[test]
fn discovers_recordable_outputs_from_all_descriptor_node_kinds() {
// `custom:` is deliberately absent: `Node` is `deny_unknown_fields` and
Expand Down
1 change: 1 addition & 0 deletions binaries/record-node/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,5 +23,6 @@ dora-recording = { workspace = true }
dora-message = { workspace = true }
eyre = { workspace = true }
aligned-vec = { version = "0.6.4", features = ["serde"] }
base64 = "0.22"
serde_json = { workspace = true }
uuid = { workspace = true, features = ["v4"] }
20 changes: 18 additions & 2 deletions binaries/record-node/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ use std::{
};

use aligned_vec::{AVec, ConstAlign};
use base64::prelude::{BASE64_STANDARD, Engine as _};
use dora_message::{
common::Timestamped,
daemon_to_daemon::InterDaemonEvent,
Expand Down Expand Up @@ -319,12 +320,27 @@ fn record_entry<W: Write>(
}
}

/// The original descriptor YAML. `dora record` passes it base64-encoded in
/// `DORA_RECORD_DESCRIPTOR_BASE64`, because `env:` values are
/// `$VAR`-expanded on the way here; the raw `DORA_RECORD_DESCRIPTOR` of an
/// older CLI is still accepted.
fn descriptor_from_env() -> eyre::Result<Vec<u8>> {
const ENCODED: &str = "DORA_RECORD_DESCRIPTOR_BASE64";
const LEGACY_RAW: &str = "DORA_RECORD_DESCRIPTOR";
match std::env::var(ENCODED) {
Ok(encoded) => BASE64_STANDARD
.decode(encoded)
.wrap_err_with(|| format!("{ENCODED} is not valid base64")),
Err(_) => Ok(std::env::var(LEGACY_RAW).unwrap_or_default().into_bytes()),
}
}

fn main() -> eyre::Result<()> {
let output_file =
std::env::var("DORA_RECORD_FILE").wrap_err("DORA_RECORD_FILE env var not set")?;
let topics_json =
std::env::var("DORA_RECORD_TOPICS").wrap_err("DORA_RECORD_TOPICS env var not set")?;
let descriptor_yaml = std::env::var("DORA_RECORD_DESCRIPTOR").unwrap_or_default();
let descriptor_yaml = descriptor_from_env()?;

// Build reverse map: input_id -> (source_node_id, source_output_id).
let reverse_map = build_reverse_map(&topics_json)?;
Expand All @@ -337,7 +353,7 @@ fn main() -> eyre::Result<()> {
version: dora_recording::FORMAT_VERSION,
start_nanos,
dataflow_id: uuid::Uuid::new_v4(),
descriptor_yaml: descriptor_yaml.into_bytes(),
descriptor_yaml,
};

let file =
Expand Down
Loading