From 67c56d9f7fa798349a4888c1566b58ffcd49e97b Mon Sep 17 00:00:00 2001 From: Arne Roomann-Kurrik Date: Thu, 3 Sep 2026 08:34:38 -0700 Subject: [PATCH 1/6] feat(services): export runtime variables from streaming targets Streaming service targets can set exports_vars to publish JSON snapshots over a private FIFO. Direct dependents consume allowlisted names through inherit_env, and exporter targets are restricted to services up. --- README.md | 48 +++- src/cli/skills.md | 7 + src/config/project.rs | 52 ++++ src/dev/daemon.rs | 4 +- src/dev/export_vars.rs | 369 +++++++++++++++++++++++++ src/dev/mod.rs | 1 + src/dev/plan.rs | 126 +++++++++ src/dev/runner.rs | 583 +++++++++++++++++++++++++++++++++++++--- src/main.rs | 58 ++++ src/plugins/mod.rs | 11 + src/targets/resolver.rs | 25 +- tests/dev_services.rs | 131 +++++++++ 12 files changed, 1382 insertions(+), 33 deletions(-) create mode 100644 src/dev/export_vars.rs diff --git a/README.md b/README.md index 8287c90..82445a1 100644 --- a/README.md +++ b/README.md @@ -376,7 +376,7 @@ If a project has a target named `services`, run it explicitly with `aster target services `. The same escape hatch works for any target name that conflicts with a built-in Aster command. -Services are mappings to ordinary `stream = true` targets. Their non-stream +Services are mappings to `stream = true` targets. Their non-stream target dependencies are pre-start steps: Aster runs them before the service starts and again before a dependency-triggered or manual restart. The same transitive target graph determines which project directories are watched. @@ -530,6 +530,52 @@ as URLs. Port references do not implicitly start services; groups remain the process-selection contract. These templates are separate from the `{files}` target capability. `open_path` controls the dashboard's browser URL. +### Runtime variable exporters + +A supervised target can publish configuration it discovers while starting or +running. Exporters are Unix-only and must be streaming service targets: + +```toml +# services/credentials/aster.toml +[targets.dev] +command = "./publish-credentials" +stream = true +exports_vars = true + +# services/api/aster.toml +[targets.dev] +command = "./start-api" +stream = true +depends_on = ["//services/credentials:dev"] + +# workspace-root aster.toml +[dev.services.credentials] +target = "//services/credentials:dev" + +[dev.services.api] +target = "//services/api:dev" +inherit_env = ["DISCOVERED_TOKEN"] +``` + +Before each exporter generation starts, Aster creates a private mode-0600 +named pipe and passes its path as `ASTER_EXPORT_VAR_PATH`. The producer writes +one complete JSON object per line; keys and values must be strings and each +record is a complete snapshot: + +```json +{"DISCOVERED_TOKEN":"example"} +``` + +Consumers wait for the first valid snapshot. Only direct target dependencies +may consume exported values, and `inherit_env` is the per-service allowlist. +Exports override env-file and ambient values; explicit service `env`, resolved +ports, Aster's internal variables, and leading command assignments remain +final. A changed effective environment restarts the affected consumer, while +identical, unused, or masked changes do not. If a producer exits or publishes +invalid data, Aster stops consumers that used its values and holds them until a +healthy snapshot is available. Records are limited to 64 KiB, malformed data +is never logged, and exporter targets cannot be run outside `aster services up`. + ### Optional service proxies A service can declare a proxy target without changing its normal launch. Pass diff --git a/src/cli/skills.md b/src/cli/skills.md index edcad69..cfbcde2 100644 --- a/src/cli/skills.md +++ b/src/cli/skills.md @@ -203,6 +203,13 @@ from their ranges as one collision-free bundle per supervisor and released at exit. Use `port_env` for direct numeric environment values and `{ports.name}` inside `env` or target commands for URLs and other composite values. +On Unix, a streaming target may set `exports_vars = true`. Aster gives each +producer generation a private `ASTER_EXPORT_VAR_PATH` FIFO. The producer writes +newline-terminated JSON objects with string keys and values. Direct target +dependencies wait for the first valid snapshot and receive only names listed in +their service's `inherit_env`; an effective value change restarts the consumer. +Exporter targets are valid only under `aster services up`. + A service may configure an optional target under `[dev.services..proxy]` with a distinct named `upstream_port` and proxy-specific `env`. `services up --proxy` keeps the service's advertised named port on the proxy and diff --git a/src/config/project.rs b/src/config/project.rs index fa8dece..3e5aa6b 100644 --- a/src/config/project.rs +++ b/src/config/project.rs @@ -118,6 +118,12 @@ pub struct RichTargetConfig { #[serde(default)] pub stream: bool, + /// Publish runtime-discovered variables to directly dependent services. + /// Exporters are supported only for streaming targets managed by the + /// service supervisor. + #[serde(default)] + pub exports_vars: bool, + /// Cache configuration overrides #[serde(default)] pub cache: Option, @@ -174,6 +180,14 @@ impl TargetConfig { } } + /// Get the runtime variable exporter flag (false for simple/alias format). + pub fn exports_vars(&self) -> bool { + match self { + TargetConfig::Simple(_) | TargetConfig::Alias(_) => false, + TargetConfig::Rich(rich) => rich.exports_vars, + } + } + /// Get cache configuration (None for simple/alias format) pub fn cache(&self) -> Option<&CacheConfig> { match self { @@ -259,6 +273,13 @@ pub(super) fn validate_aster_config( continue; }; + if target.exports_vars && !target.stream { + anyhow::bail!( + "Target '{target_name}' in {} sets exports_vars = true but is not streaming; variable exporters must also set stream = true", + path.display() + ); + } + for capability in &target.capabilities { if capability != "files_list" { anyhow::bail!( @@ -591,6 +612,37 @@ command = "pytest" assert!(!test.stream()); } + #[test] + fn test_variable_exporter_requires_streaming_target() { + let tmp = tempfile::tempdir().unwrap(); + let toml_path = tmp.path().join("aster.toml"); + std::fs::write( + &toml_path, + r#" +[targets.exporter] +command = "./publish" +exports_vars = true +"#, + ) + .unwrap(); + + let error = parse_aster_toml(&toml_path).unwrap_err().to_string(); + assert!(error.contains("exports_vars = true")); + + std::fs::write( + &toml_path, + r#" +[targets.exporter] +command = "./publish" +stream = true +exports_vars = true +"#, + ) + .unwrap(); + let config = parse_aster_toml(&toml_path).unwrap(); + assert!(config.targets["exporter"].exports_vars()); + } + #[test] fn test_parse_alias_bare_name() { let tmp = tempfile::tempdir().unwrap(); diff --git a/src/dev/daemon.rs b/src/dev/daemon.rs index 13be630..86a6dfd 100644 --- a/src/dev/daemon.rs +++ b/src/dev/daemon.rs @@ -240,7 +240,9 @@ mod platform { const PID_NAME: &str = "daemon.pid"; const LOG_NAME: &str = "daemon.log"; const READY_NAME: &str = "ready.sock"; - const START_TIMEOUT: Duration = Duration::from_secs(20); + // Exporting services own a 20-second first-snapshot deadline. Leave a + // bounded handoff margin for the supervisor to report that result. + const START_TIMEOUT: Duration = Duration::from_secs(25); const CONNECT_TIMEOUT: Duration = Duration::from_secs(10); const IDLE_GRACE: Duration = Duration::from_millis(500); const STOP_GRACE: Duration = Duration::from_secs(5); diff --git a/src/dev/export_vars.rs b/src/dev/export_vars.rs new file mode 100644 index 0000000..02419c4 --- /dev/null +++ b/src/dev/export_vars.rs @@ -0,0 +1,369 @@ +//! Private JSONL endpoints used by supervised variable-exporting targets. + +use std::collections::HashMap; + +pub(crate) const EXPORT_PATH_ENV: &str = "ASTER_EXPORT_VAR_PATH"; +const MAX_RECORD_BYTES: usize = 64 * 1024; + +#[derive(Debug)] +pub(crate) enum ExportEvent { + Snapshot { + producer: String, + generation: u64, + values: HashMap, + }, + Invalid { + producer: String, + generation: u64, + reason: String, + }, +} + +#[derive(Default)] +struct JsonlParser { + buffer: Vec, + discarding_oversized: bool, +} + +impl JsonlParser { + fn push(&mut self, mut bytes: &[u8]) -> Vec, String>> { + let mut records = Vec::new(); + if self.discarding_oversized { + let Some(newline) = bytes.iter().position(|byte| *byte == b'\n') else { + return records; + }; + bytes = &bytes[newline + 1..]; + self.discarding_oversized = false; + } + self.buffer.extend_from_slice(bytes); + + loop { + let Some(newline) = self.buffer.iter().position(|byte| *byte == b'\n') else { + if self.buffer.len() > MAX_RECORD_BYTES { + self.buffer.clear(); + self.discarding_oversized = true; + records.push(Err("record exceeds 64 KiB".to_string())); + } + break; + }; + let mut record = self.buffer.drain(..=newline).collect::>(); + record.pop(); + if record.len() > MAX_RECORD_BYTES { + records.push(Err("record exceeds 64 KiB".to_string())); + continue; + } + records.push(parse_snapshot(&record)); + } + records + } +} + +fn parse_snapshot(bytes: &[u8]) -> Result, String> { + use serde::de::{Error as _, MapAccess, Visitor}; + + struct SnapshotVisitor; + impl<'de> Visitor<'de> for SnapshotVisitor { + type Value = HashMap; + + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("a JSON object with unique environment-name keys and string values") + } + + fn visit_map(self, mut map: A) -> Result + where + A: MapAccess<'de>, + { + let mut values = HashMap::new(); + while let Some(key) = map.next_key::()? { + let value = map.next_value::()?; + if values.insert(key.clone(), value).is_some() { + return Err(A::Error::custom(format!("duplicate key '{key}'"))); + } + } + Ok(values) + } + } + + let mut deserializer = serde_json::Deserializer::from_slice(bytes); + let values = serde::de::Deserializer::deserialize_map(&mut deserializer, SnapshotVisitor) + .map_err(|error| format!("invalid JSON snapshot: {error}"))?; + deserializer + .end() + .map_err(|error| format!("trailing JSON data: {error}"))?; + for (name, value) in &values { + if !valid_environment_name(name) { + return Err(format!("invalid environment variable name '{name}'")); + } + if value.contains('\0') { + return Err(format!("environment variable '{name}' contains NUL")); + } + } + Ok(values) +} + +fn valid_environment_name(name: &str) -> bool { + let mut bytes = name.bytes(); + matches!(bytes.next(), Some(b'A'..=b'Z' | b'a'..=b'z' | b'_')) + && bytes.all(|byte| byte.is_ascii_alphanumeric() || byte == b'_') +} + +#[cfg(unix)] +mod platform { + use super::*; + use anyhow::{Context, Result}; + use std::ffi::CString; + use std::fs::{self, File, OpenOptions}; + use std::io::Read; + use std::os::unix::ffi::OsStrExt; + use std::os::unix::fs::{ + DirBuilderExt, FileTypeExt, MetadataExt, OpenOptionsExt, PermissionsExt, + }; + use std::path::{Path, PathBuf}; + use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::mpsc::Sender; + use std::sync::Arc; + use std::thread::JoinHandle; + use std::time::Duration; + + pub(crate) struct ExportEndpoint { + directory: PathBuf, + path: PathBuf, + _keepalive_writer: File, + stop: Arc, + reader: Option>, + } + + impl ExportEndpoint { + pub(crate) fn create( + producer: &str, + generation: u64, + events: Sender, + ) -> Result { + cleanup_stale_endpoints(); + let mut random = [0_u8; 16]; + getrandom::fill(&mut random).map_err(|error| { + anyhow::anyhow!("failed to generate variable endpoint name: {error}") + })?; + let nonce = random + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + let directory = std::env::temp_dir().join(format!( + "aster-export-vars-{}-{generation}-{nonce}", + std::process::id() + )); + let mut directory_builder = fs::DirBuilder::new(); + directory_builder.mode(0o700); + directory_builder.create(&directory).with_context(|| { + format!("failed to create variable endpoint {}", directory.display()) + })?; + fs::set_permissions(&directory, fs::Permissions::from_mode(0o700))?; + let path = directory.join("vars.jsonl"); + make_fifo(&path)?; + fs::set_permissions(&path, fs::Permissions::from_mode(0o600))?; + + let reader = OpenOptions::new() + .read(true) + .custom_flags(libc::O_NONBLOCK | libc::O_CLOEXEC) + .open(&path) + .with_context(|| format!("failed to open variable reader {}", path.display()))?; + let keepalive_writer = OpenOptions::new() + .write(true) + .custom_flags(libc::O_NONBLOCK | libc::O_CLOEXEC) + .open(&path) + .with_context(|| format!("failed to open variable keepalive {}", path.display()))?; + let stop = Arc::new(AtomicBool::new(false)); + let reader_stop = stop.clone(); + let event_producer = producer.to_string(); + let handle = std::thread::spawn(move || { + read_records(reader, &event_producer, generation, &events, &reader_stop); + }); + Ok(Self { + directory, + path, + _keepalive_writer: keepalive_writer, + stop, + reader: Some(handle), + }) + } + + pub(crate) fn path(&self) -> &Path { + &self.path + } + } + + impl Drop for ExportEndpoint { + fn drop(&mut self) { + self.stop.store(true, Ordering::SeqCst); + if let Some(reader) = self.reader.take() { + let _ = reader.join(); + } + let _ = fs::remove_file(&self.path); + let _ = fs::remove_dir(&self.directory); + } + } + + fn make_fifo(path: &Path) -> Result<()> { + let path = CString::new(path.as_os_str().as_bytes()) + .context("variable endpoint path contains NUL")?; + let result = unsafe { libc::mkfifo(path.as_ptr(), 0o600) }; + if result != 0 { + return Err(std::io::Error::last_os_error()).context("failed to create variable FIFO"); + } + Ok(()) + } + + fn cleanup_stale_endpoints() { + let temp = std::env::temp_dir(); + let Ok(entries) = fs::read_dir(&temp) else { + return; + }; + let current_uid = unsafe { libc::geteuid() }; + for entry in entries.flatten() { + let name = entry.file_name(); + let Some(name) = name.to_str() else { continue }; + let Some(rest) = name.strip_prefix("aster-export-vars-") else { + continue; + }; + let Some(pid) = rest + .split('-') + .next() + .and_then(|value| value.parse::().ok()) + else { + continue; + }; + if process_exists(pid) { + continue; + } + let path = entry.path(); + let Ok(metadata) = fs::symlink_metadata(&path) else { + continue; + }; + if !metadata.file_type().is_dir() || metadata.uid() != current_uid { + continue; + } + let fifo = path.join("vars.jsonl"); + if fs::symlink_metadata(&fifo).is_ok_and(|metadata| { + metadata.file_type().is_fifo() && metadata.uid() == current_uid + }) { + let _ = fs::remove_file(&fifo); + } + let _ = fs::remove_dir(&path); + } + } + + fn process_exists(pid: i32) -> bool { + let result = unsafe { libc::kill(pid, 0) }; + result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM) + } + + fn read_records( + mut reader: File, + producer: &str, + generation: u64, + events: &Sender, + stop: &AtomicBool, + ) { + let mut parser = JsonlParser::default(); + let mut buffer = [0_u8; 8192]; + while !stop.load(Ordering::SeqCst) { + match reader.read(&mut buffer) { + Ok(0) => std::thread::sleep(Duration::from_millis(10)), + Ok(size) => { + for record in parser.push(&buffer[..size]) { + let event = match record { + Ok(values) => ExportEvent::Snapshot { + producer: producer.to_string(), + generation, + values, + }, + Err(reason) => ExportEvent::Invalid { + producer: producer.to_string(), + generation, + reason, + }, + }; + if events.send(event).is_err() { + return; + } + } + } + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + std::thread::sleep(Duration::from_millis(10)); + } + Err(_) => return, + } + } + } +} + +#[cfg(unix)] +pub(crate) use platform::ExportEndpoint; + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parser_handles_split_and_combined_records() { + let mut parser = JsonlParser::default(); + assert!(parser.push(br#"{"ONE":"1""#).is_empty()); + let records = parser.push(b"}\n{\"TWO\":\"2\"}\n"); + assert_eq!(records.len(), 2); + assert_eq!(records[0].as_ref().unwrap()["ONE"], "1"); + assert_eq!(records[1].as_ref().unwrap()["TWO"], "2"); + } + + #[test] + fn parser_recovers_after_oversized_record() { + let mut parser = JsonlParser::default(); + let oversized = vec![b'x'; MAX_RECORD_BYTES + 1]; + assert!(parser.push(&oversized).pop().unwrap().is_err()); + let records = parser.push(b"discarded\n{\"OK\":\"yes\"}\n"); + assert_eq!(records.len(), 1); + assert_eq!(records[0].as_ref().unwrap()["OK"], "yes"); + } + + #[test] + fn parser_rejects_duplicates_invalid_names_and_non_strings() { + for input in [ + b"{\"A\":\"1\",\"A\":\"2\"}\n".as_slice(), + b"{\"BAD-NAME\":\"1\"}\n".as_slice(), + b"{\"A\":1}\n".as_slice(), + ] { + assert!(JsonlParser::default().push(input).pop().unwrap().is_err()); + } + } + + #[cfg(unix)] + #[test] + fn endpoint_is_private_and_publishes_without_logging_values() { + use std::io::Write; + use std::os::unix::fs::{FileTypeExt, PermissionsExt}; + use std::sync::mpsc; + use std::time::Duration; + + let (tx, rx) = mpsc::channel(); + let endpoint = ExportEndpoint::create("producer", 7, tx).unwrap(); + let metadata = std::fs::metadata(endpoint.path()).unwrap(); + assert!(metadata.file_type().is_fifo()); + assert_eq!(metadata.permissions().mode() & 0o777, 0o600); + let mut writer = std::fs::OpenOptions::new() + .write(true) + .open(endpoint.path()) + .unwrap(); + writer.write_all(b"{\"TOKEN\":\"secret\"}\n").unwrap(); + match rx.recv_timeout(Duration::from_secs(1)).unwrap() { + ExportEvent::Snapshot { + producer, + generation, + values, + } => { + assert_eq!(producer, "producer"); + assert_eq!(generation, 7); + assert_eq!(values["TOKEN"], "secret"); + } + event => panic!("unexpected event: {event:?}"), + } + } +} diff --git a/src/dev/mod.rs b/src/dev/mod.rs index 6b447bc..5d11013 100644 --- a/src/dev/mod.rs +++ b/src/dev/mod.rs @@ -2,6 +2,7 @@ mod daemon; mod dashboard; +mod export_vars; mod log_files; mod plan; mod port_allocator; diff --git a/src/dev/plan.rs b/src/dev/plan.rs index ea61402..b36864b 100644 --- a/src/dev/plan.rs +++ b/src/dev/plan.rs @@ -6,6 +6,7 @@ use anyhow::{anyhow, bail, Context, Result}; use crate::config::{DevPortConfig, DevWorkspaceConfig}; use crate::discovery::DiscoveredProject; +use crate::executor::command::parse_command; use crate::graph::TargetGraph; use crate::plugins::{PluginRegistry, Target}; use crate::watch::WatchPlan; @@ -29,6 +30,12 @@ pub struct ServicePlan { pub port: Option, pub open_url: Option, pub env: HashMap, + /// Service names whose streaming targets must be ready first. + pub dependencies: Vec, + /// Exported variable names this service is permitted to consume. + pub inherit_env: HashSet, + /// Higher-precedence environment names that exported snapshots cannot replace. + pub protected_env: HashSet, pub watch: WatchPlan, } @@ -112,6 +119,11 @@ pub fn resolve_dev_plan( .copied(); let service_file_env = load_env_files(workspace_root, &service.env_files)?; let mut service_env = service_file_env; + let mut protected_env = service.env.keys().cloned().collect::>(); + protected_env.extend(service.port_env.keys().cloned()); + protected_env.insert("ASTER_SERVICE_NAME".to_string()); + protected_env.insert("ASTER_SERVICE_PORT".to_string()); + protected_env.insert("ASTER_RESOLVED_PORTS".to_string()); for key in &service.inherit_env { validate_environment_key(name, key)?; if let Ok(value) = env::var(key) { @@ -215,6 +227,13 @@ pub fn resolve_dev_plan( }; target.command = expand_template(&target.command, port, &ports) .with_context(|| format!("invalid command for service '{name}'"))?; + protected_env.extend( + parse_command(&target.command)? + .env + .into_iter() + .map(|(name, _)| name), + ); + protected_env.insert("ASTER_EXPORT_VAR_PATH".to_string()); services.push(ServicePlan { name: name.clone(), @@ -224,6 +243,9 @@ pub fn resolve_dev_plan( port, open_url, env: service_env, + dependencies: Vec::new(), + inherit_env: service.inherit_env.iter().cloned().collect(), + protected_env, watch, }); @@ -243,6 +265,20 @@ pub fn resolve_dev_plan( target.command = expand_proxy_template(&target.command, listen_port, upstream_port, &ports) .with_context(|| format!("invalid command for service '{proxy_name}'"))?; + let mut proxy_protected_env = HashSet::from([ + "ASTER_SERVICE_NAME".to_string(), + "ASTER_SERVICE_PORT".to_string(), + "ASTER_PROXY_SERVICE_NAME".to_string(), + "ASTER_PROXY_LISTEN_PORT".to_string(), + "ASTER_PROXY_UPSTREAM_PORT".to_string(), + "ASTER_EXPORT_VAR_PATH".to_string(), + ]); + proxy_protected_env.extend( + parse_command(&target.command)? + .env + .into_iter() + .map(|(name, _)| name), + ); let mut proxy_env = HashMap::new(); for (key, value) in &proxy.env { proxy_env.insert( @@ -272,6 +308,9 @@ pub fn resolve_dev_plan( port: Some(listen_port), open_url: None, env: proxy_env, + dependencies: vec![name.clone()], + inherit_env: HashSet::new(), + protected_env: proxy_protected_env, watch, }); } @@ -284,6 +323,8 @@ pub fn resolve_dev_plan( } } + let services = order_service_graph(services, projects)?; + Ok(DevPlan { services, ports, @@ -292,6 +333,91 @@ pub fn resolve_dev_plan( }) } +/// Connect selected service instances through direct streaming-target edges and +/// return a stable topological order. Existing configured order remains the tie +/// breaker because `services` arrives in that order. +fn order_service_graph( + mut services: Vec, + projects: &[DiscoveredProject], +) -> Result> { + let exporter_addresses = projects + .iter() + .flat_map(|project| { + project + .targets + .iter() + .filter(|(_, target)| target.exports_vars()) + .map(move |(name, _)| format!("//{}:{name}", project.relative_path.display())) + }) + .collect::>(); + #[cfg(not(unix))] + if !exporter_addresses.is_empty() { + bail!("exports_vars targets require the Unix service supervisor"); + } + let mut instances: HashMap> = HashMap::new(); + for service in &services { + instances + .entry(service.target_address.clone()) + .or_default() + .push(service.name.clone()); + } + for address in &exporter_addresses { + if let Some(names) = instances.get(address) { + if names.len() != 1 { + bail!( + "variable exporter {address} has multiple selected service instances: {}", + names.join(", ") + ); + } + } + } + + for service in &mut services { + for dependency in &service.target.depends_on { + match instances.get(dependency) { + Some(names) + if names.len() == 1 && !service.dependencies.contains(&names[0]) => + { + service.dependencies.push(names[0].clone()); + } + Some(names) if names.len() == 1 => {} + Some(names) if exporter_addresses.contains(dependency) => bail!( + "variable exporter {dependency} has multiple selected service instances: {}", + names.join(", ") + ), + Some(_) => {} + None if exporter_addresses.contains(dependency) => bail!( + "service '{}' depends on variable exporter {dependency}, but that producer is not selected", + service.name + ), + None => {} + } + } + } + + let mut ordered = Vec::with_capacity(services.len()); + let mut emitted = HashSet::new(); + while !services.is_empty() { + let Some(index) = services.iter().position(|service| { + service + .dependencies + .iter() + .all(|dependency| emitted.contains(dependency)) + }) else { + let names = services + .iter() + .map(|service| service.name.as_str()) + .collect::>() + .join(", "); + bail!("service dependency cycle among: {names}"); + }; + let service = services.remove(index); + emitted.insert(service.name.clone()); + ordered.push(service); + } + Ok(ordered) +} + fn resolve_stream_target( service_name: &str, target_address: &str, diff --git a/src/dev/runner.rs b/src/dev/runner.rs index 55ac084..b3c09ea 100644 --- a/src/dev/runner.rs +++ b/src/dev/runner.rs @@ -20,6 +20,8 @@ use crate::graph::TargetGraph; use crate::watch::WorkspaceIgnore; use super::dashboard::{Dashboard, DashboardAction, ServiceState, TerminalGuard}; +#[cfg(unix)] +use super::export_vars::{ExportEndpoint, ExportEvent, EXPORT_PATH_ENV}; use super::log_files::ServiceLogFiles; use super::plan::{DevPlan, ServicePlan}; use super::process::{LogEvent, ProcessLogSenders, ServiceProcess}; @@ -35,16 +37,43 @@ pub struct DevOptions { struct Runtime { process: Option, + #[cfg(unix)] + export_endpoint: Option, + generation: u64, + applied_exports: HashMap, + applied_epochs: HashMap, + applied_sources: HashSet, +} + +#[derive(Default)] +struct ProducerState { + generation: u64, + epoch: u64, + values: Option>, + healthy: bool, + waiting_since: Option, +} + +#[derive(Clone, Debug, Default)] +struct ResolvedExports { + values: HashMap, + epochs: HashMap, + sources: HashSet, } enum StartOutcome { - Running(ServiceProcess), + Running { + process: ServiceProcess, + #[cfg(unix)] + export_endpoint: Option, + }, Stopped, } struct StartResult { service: String, outcome: StartOutcome, + resolved_exports: ResolvedExports, } struct ActiveStart { @@ -100,23 +129,34 @@ pub fn run_dev( } else { (None, None) }; - super::daemon::register_supervisor_ready( - workspace_root, - plan.services - .iter() - .map(|service| service.name.clone()) - .collect(), - plan.ports - .iter() - .map(|(name, port)| (name.clone(), *port)) - .collect(), - )?; + let mut supervisor_registered = false; let mut dashboard = Dashboard::new(&plan.services); let mut runtimes: HashMap = plan .services .iter() - .map(|service| (service.name.clone(), Runtime { process: None })) + .map(|service| { + ( + service.name.clone(), + Runtime { + process: None, + #[cfg(unix)] + export_endpoint: None, + generation: 0, + applied_exports: HashMap::new(), + applied_epochs: HashMap::new(), + applied_sources: HashSet::new(), + }, + ) + }) .collect(); + let mut producers = plan + .services + .iter() + .filter(|service| service.target.exports_vars()) + .map(|service| (service.name.clone(), ProducerState::default())) + .collect::>(); + #[cfg(unix)] + let (export_tx, export_rx) = mpsc::channel::(); let projects = Arc::new(projects); let (start_tx, start_rx) = mpsc::channel::(); let mut active_start: Option = None; @@ -139,6 +179,11 @@ pub fn run_dev( let mut watch_deadline: Option = None; let mut quitting = false; + if producers.is_empty() { + register_ready(workspace_root, &plan)?; + supervisor_registered = true; + } + let run_result = (|| -> Result<()> { while !quitting && !shutdown.load(Ordering::SeqCst) @@ -167,21 +212,64 @@ pub fn run_dev( ); } } - let runtime = runtimes - .get_mut(&result.service) - .expect("start result service exists"); match result.outcome { - StartOutcome::Running(process) => { - runtime.process = Some(process); - set_service_state( - &mut dashboard, + StartOutcome::Running { + mut process, + #[cfg(unix)] + export_endpoint, + } => { + let service = plan + .services + .iter() + .find(|service| service.name == result.service) + .expect("start result service exists"); + let current = resolve_exported_environment(service, &producers)?; + if current.as_ref().map(|resolved| &resolved.values) + != Some(&result.resolved_exports.values) + { + process.terminate(Duration::from_secs(3)); #[cfg(unix)] - session.as_ref(), - &result.service, - ServiceState::Running, - ); + drop(export_endpoint); + queue_restart( + &mut pending_starts, + &result.service, + "exported variables changed during start", + &system_tx, + &mut dashboard, + ); + } else { + let runtime = runtimes + .get_mut(&result.service) + .expect("start result service exists"); + runtime.process = Some(process); + #[cfg(unix)] + { + runtime.export_endpoint = export_endpoint; + } + runtime.applied_exports = result.resolved_exports.values; + runtime.applied_sources = result.resolved_exports.sources; + runtime.applied_epochs = current + .map(|resolved| resolved.epochs) + .unwrap_or(result.resolved_exports.epochs); + set_service_state( + &mut dashboard, + #[cfg(unix)] + session.as_ref(), + &result.service, + ServiceState::Running, + ); + } } StartOutcome::Stopped => { + if producers + .get(&result.service) + .is_some_and(|producer| producer.values.is_none()) + { + anyhow::bail!( + "variable exporter '{}' exited before publishing its initial snapshot", + result.service + ); + } set_service_state( &mut dashboard, #[cfg(unix)] @@ -194,8 +282,144 @@ pub fn run_dev( suppress_until = Instant::now() + Duration::from_millis(700); needs_draw = true; } + + #[cfg(unix)] + while let Ok(event) = export_rx.try_recv() { + match event { + ExportEvent::Snapshot { + producer, + generation, + values, + } => { + let Some(state) = producers.get_mut(&producer) else { + continue; + }; + if state.generation != generation + || runtimes[&producer].generation != generation + { + continue; + } + state.epoch = state.epoch.saturating_add(1); + state.values = Some(values); + state.healthy = true; + state.waiting_since = None; + emit_system( + &system_tx, + &producer, + format!("accepted exported variable snapshot epoch {}", state.epoch), + false, + ); + for service in &plan.services { + if service.name == producer + || !service.dependencies.contains(&producer) + || runtimes[&service.name].process.is_none() + { + continue; + } + let desired = resolve_exported_environment(service, &producers)?; + if desired.as_ref().map(|resolved| &resolved.values) + != Some(&runtimes[&service.name].applied_exports) + { + queue_restart( + &mut pending_starts, + &service.name, + "effective exported environment changed", + &system_tx, + &mut dashboard, + ); + } else if let Some(resolved) = desired { + let runtime = + runtimes.get_mut(&service.name).expect("runtime exists"); + runtime.applied_epochs = resolved.epochs; + runtime.applied_sources = resolved.sources; + } + } + } + ExportEvent::Invalid { + producer, + generation, + reason, + } => { + let Some(state) = producers.get_mut(&producer) else { + continue; + }; + if state.generation != generation { + continue; + } + state.epoch = state.epoch.saturating_add(1); + state.healthy = false; + emit_system( + &system_tx, + &producer, + format!("rejected exported variable update: {reason}"), + true, + ); + stop_dependent_services( + &producer, + &plan, + &mut runtimes, + &mut pending_starts, + &system_tx, + &mut dashboard, + #[cfg(unix)] + session.as_ref(), + ); + } + } + needs_draw = true; + } + + if !supervisor_registered + && producers + .values() + .all(|producer| producer.healthy && producer.values.is_some()) + { + register_ready(workspace_root, &plan)?; + supervisor_registered = true; + } + + if let Some((name, _)) = producers.iter().find(|(_, producer)| { + producer.waiting_since.is_some_and(|started| { + started.elapsed() >= Duration::from_secs(20) && producer.values.is_none() + }) + }) { + anyhow::bail!( + "variable exporter '{name}' did not publish an initial snapshot within 20 seconds" + ); + } if active_start.is_none() { - if let Some((name, reason)) = pending_starts.pop_front() { + let mut startable = None; + for (index, (name, _)) in pending_starts.iter().enumerate() { + let service = plan + .services + .iter() + .find(|service| service.name == *name) + .expect("queued service exists"); + if service_is_startable(service, &runtimes, &producers)? { + startable = Some(index); + break; + } + } + if let Some(index) = startable { + let (name, reason) = pending_starts + .remove(index) + .expect("startable queue position exists"); + if reason != "initial start" + && producers + .get(&name) + .is_some_and(|producer| producer.values.is_some()) + { + stop_dependent_services( + &name, + &plan, + &mut runtimes, + &mut pending_starts, + &system_tx, + &mut dashboard, + #[cfg(unix)] + session.as_ref(), + ); + } let service = plan .services .iter() @@ -204,6 +428,8 @@ pub fn run_dev( let before_start = options .watch .then(|| suppressed_path_snapshot(&plan, workspace_root, &ignore)); + let resolved_exports = + resolve_exported_environment(service, &producers)?.unwrap_or_default(); let handle = begin_start_service( service, workspace_root, @@ -219,6 +445,10 @@ pub fn run_dev( session.as_ref(), runtimes.get_mut(&name).expect("runtime exists"), &reason, + resolved_exports, + #[cfg(unix)] + export_tx.clone(), + producers.get_mut(&name), start_tx.clone(), ); active_start = Some(ActiveStart { @@ -231,13 +461,17 @@ pub fn run_dev( } for service in &plan.services { - let runtime = runtimes.get_mut(&service.name).expect("runtime exists"); - let exited = runtime + let exited = runtimes + .get_mut(&service.name) + .expect("runtime exists") .process .as_mut() .and_then(|process| process.poll().ok().flatten()); if let Some(code) = exited { + let runtime = runtimes.get_mut(&service.name).expect("runtime exists"); runtime.process.take(); + #[cfg(unix)] + runtime.export_endpoint.take(); set_service_state( &mut dashboard, #[cfg(unix)] @@ -251,6 +485,29 @@ pub fn run_dev( format!("process exited with code {code}"), true, ); + if let Some(producer) = producers.get_mut(&service.name) { + let had_snapshot = producer.values.is_some(); + producer.epoch = producer.epoch.saturating_add(1); + producer.values = None; + producer.healthy = false; + producer.waiting_since = None; + if !had_snapshot && !supervisor_registered { + anyhow::bail!( + "variable exporter '{}' exited before publishing its initial snapshot", + service.name + ); + } + stop_dependent_services( + &service.name, + &plan, + &mut runtimes, + &mut pending_starts, + &system_tx, + &mut dashboard, + #[cfg(unix)] + session.as_ref(), + ); + } needs_draw = true; } } @@ -511,11 +768,22 @@ pub fn run_dev( } let mut shutdown_processes = Vec::new(); while let Ok(result) = start_rx.try_recv() { - if let StartOutcome::Running(process) = result.outcome { + if let StartOutcome::Running { + process, + #[cfg(unix)] + export_endpoint, + } = result.outcome + { + #[cfg(unix)] + drop(export_endpoint); shutdown_processes.push(process); } } for service in plan.services.iter().rev() { + #[cfg(unix)] + runtimes + .get_mut(&service.name) + .and_then(|runtime| runtime.export_endpoint.take()); if let Some(process) = runtimes .get_mut(&service.name) .and_then(|runtime| runtime.process.take()) @@ -546,6 +814,140 @@ pub fn run_dev( Ok(()) } +fn register_ready(workspace_root: &Path, plan: &DevPlan) -> Result<()> { + super::daemon::register_supervisor_ready( + workspace_root, + plan.services + .iter() + .map(|service| service.name.clone()) + .collect(), + plan.ports + .iter() + .map(|(name, port)| (name.clone(), *port)) + .collect(), + )?; + Ok(()) +} + +fn resolve_exported_environment( + service: &ServicePlan, + producers: &HashMap, +) -> Result> { + let mut environment = HashMap::new(); + let mut epochs = HashMap::new(); + let mut sources = HashMap::::new(); + let mut used_sources = HashSet::new(); + for dependency in &service.dependencies { + let Some(producer) = producers.get(dependency) else { + continue; + }; + if !producer.healthy { + return Ok(None); + } + let Some(values) = producer.values.as_ref() else { + return Ok(None); + }; + for name in &service.inherit_env { + if service.protected_env.contains(name) { + continue; + } + let Some(value) = values.get(name) else { + continue; + }; + if let Some(existing) = sources.insert(name.clone(), dependency.clone()) { + anyhow::bail!( + "service '{}' consumes exported variable '{name}' from both '{existing}' and '{dependency}'", + service.name + ); + } + environment.insert(name.clone(), value.clone()); + used_sources.insert(dependency.clone()); + } + epochs.insert(dependency.clone(), producer.epoch); + } + + Ok(Some(ResolvedExports { + values: environment, + epochs, + sources: used_sources, + })) +} + +fn service_is_startable( + service: &ServicePlan, + runtimes: &HashMap, + producers: &HashMap, +) -> Result { + if service.dependencies.iter().any(|dependency| { + runtimes + .get(dependency) + .is_none_or(|runtime| runtime.process.is_none()) + }) { + return Ok(false); + } + Ok(resolve_exported_environment(service, producers)?.is_some()) +} + +#[allow(clippy::too_many_arguments)] +fn stop_dependent_services( + producer: &str, + plan: &DevPlan, + runtimes: &mut HashMap, + pending: &mut VecDeque<(String, String)>, + system_tx: &std::sync::mpsc::Sender, + dashboard: &mut Dashboard, + #[cfg(unix)] session: Option<&SessionServer>, +) { + let mut affected = HashSet::new(); + let mut frontier = vec![producer.to_string()]; + while let Some(current) = frontier.pop() { + for service in &plan.services { + if affected.contains(&service.name) + || !runtimes[&service.name].applied_sources.contains(¤t) + { + continue; + } + affected.insert(service.name.clone()); + if service.target.exports_vars() { + frontier.push(service.name.clone()); + } + } + } + + for service in plan.services.iter().rev() { + if !affected.contains(&service.name) { + continue; + } + let runtime = runtimes.get_mut(&service.name).expect("runtime exists"); + #[cfg(unix)] + runtime.export_endpoint.take(); + if let Some(mut process) = runtime.process.take() { + process.terminate(Duration::from_secs(3)); + } + runtime.applied_exports.clear(); + runtime.applied_epochs.clear(); + runtime.applied_sources.clear(); + set_service_state( + dashboard, + #[cfg(unix)] + session, + &service.name, + ServiceState::Restarting, + ); + } + for service in &plan.services { + if affected.contains(&service.name) { + queue_restart( + pending, + &service.name, + "required exported variables unavailable", + system_tx, + dashboard, + ); + } + } +} + fn set_service_state( dashboard: &mut Dashboard, #[cfg(unix)] session: Option<&SessionServer>, @@ -665,11 +1067,25 @@ fn begin_start_service( #[cfg(unix)] session: Option<&SessionServer>, runtime: &mut Runtime, reason: &str, + resolved_exports: ResolvedExports, + #[cfg(unix)] export_tx: std::sync::mpsc::Sender, + producer_state: Option<&mut ProducerState>, result_tx: std::sync::mpsc::Sender, ) -> std::thread::JoinHandle<()> { + #[cfg(unix)] + runtime.export_endpoint.take(); if let Some(mut process) = runtime.process.take() { process.terminate(Duration::from_secs(3)); } + runtime.generation = runtime.generation.saturating_add(1); + let generation = runtime.generation; + if let Some(producer) = producer_state { + producer.generation = generation; + producer.epoch = producer.epoch.saturating_add(1); + producer.values = None; + producer.healthy = false; + producer.waiting_since = Some(Instant::now()); + } let state = if reason == "initial start" { ServiceState::Starting } else { @@ -707,7 +1123,8 @@ fn begin_start_service( let target_address = service.target_address.clone(); let target = service.target.clone(); let project_root = service.project_root.clone(); - let env = service.env.clone(); + let mut env = service.env.clone(); + env.extend(resolved_exports.values.clone()); let workspace_root = workspace_root.to_path_buf(); let log_tx = log_tx.clone(); let system_tx = system_tx.clone(); @@ -735,6 +1152,7 @@ fn begin_start_service( let _ = result_tx.send(StartResult { service: service_name, outcome: StartOutcome::Stopped, + resolved_exports, }); return; } @@ -748,6 +1166,7 @@ fn begin_start_service( let _ = result_tx.send(StartResult { service: service_name, outcome: StartOutcome::Stopped, + resolved_exports, }); return; } @@ -757,6 +1176,34 @@ fn begin_start_service( format!("starting {target_address}"), false, ); + #[cfg(unix)] + let export_endpoint = if target.exports_vars() { + match ExportEndpoint::create(&service_name, generation, export_tx) { + Ok(endpoint) => { + env.insert( + EXPORT_PATH_ENV.to_string(), + endpoint.path().to_string_lossy().into_owned(), + ); + Some(endpoint) + } + Err(error) => { + emit_system( + &system_tx, + &service_name, + format!("failed to create variable export endpoint: {error:#}"), + true, + ); + let _ = result_tx.send(StartResult { + service: service_name, + outcome: StartOutcome::Stopped, + resolved_exports, + }); + return; + } + } + } else { + None + }; let outcome = match ServiceProcess::spawn( &service_name, &target, @@ -764,7 +1211,11 @@ fn begin_start_service( &env, ProcessLogSenders::new(log_tx, system_tx.clone(), durable_log_tx, ui), ) { - Ok(process) => StartOutcome::Running(process), + Ok(process) => StartOutcome::Running { + process, + #[cfg(unix)] + export_endpoint, + }, Err(error) => { emit_system( &system_tx, @@ -778,6 +1229,7 @@ fn begin_start_service( let _ = result_tx.send(StartResult { service: service_name, outcome, + resolved_exports, }); }) } @@ -1048,6 +1500,77 @@ fn is_meaningful_event(kind: ¬ify::EventKind) -> bool { mod tests { use super::*; + fn export_consumer(dependencies: &[&str], inherit: &[&str], protected: &[&str]) -> ServicePlan { + ServicePlan { + name: "consumer".to_string(), + target_address: "//consumer:dev".to_string(), + target: crate::plugins::Target::default(), + project_root: PathBuf::new(), + port: None, + open_url: None, + env: HashMap::new(), + dependencies: dependencies + .iter() + .map(|value| (*value).to_string()) + .collect(), + inherit_env: inherit.iter().map(|value| (*value).to_string()).collect(), + protected_env: protected.iter().map(|value| (*value).to_string()).collect(), + watch: crate::watch::WatchPlan::empty(), + } + } + + fn ready_producer(epoch: u64, values: &[(&str, &str)]) -> ProducerState { + ProducerState { + generation: 1, + epoch, + values: Some( + values + .iter() + .map(|(name, value)| ((*name).to_string(), (*value).to_string())) + .collect(), + ), + healthy: true, + waiting_since: None, + } + } + + #[test] + fn exported_environment_applies_allowlist_and_precedence() { + let service = export_consumer(&["producer"], &["TOKEN", "PRIVATE"], &["TOKEN"]); + let producers = HashMap::from([( + "producer".to_string(), + ready_producer(4, &[("TOKEN", "masked"), ("PRIVATE", "selected")]), + )]); + + let resolved = resolve_exported_environment(&service, &producers) + .unwrap() + .unwrap(); + assert_eq!( + resolved.values, + HashMap::from([("PRIVATE".to_string(), "selected".to_string())]) + ); + assert_eq!(resolved.epochs["producer"], 4); + assert_eq!(resolved.sources, HashSet::from(["producer".to_string()])); + } + + #[test] + fn exported_environment_waits_for_health_and_rejects_collisions() { + let service = export_consumer(&["first", "second"], &["TOKEN"], &[]); + let mut producers = HashMap::from([ + ("first".to_string(), ready_producer(1, &[("TOKEN", "one")])), + ("second".to_string(), ready_producer(1, &[("TOKEN", "two")])), + ]); + let error = resolve_exported_environment(&service, &producers) + .unwrap_err() + .to_string(); + assert!(error.contains("from both 'first' and 'second'")); + + producers.get_mut("second").unwrap().healthy = false; + assert!(resolve_exported_environment(&service, &producers) + .unwrap() + .is_none()); + } + #[test] fn delayed_generated_event_is_ignored_but_changed_identity_restarts() { let temp = tempfile::tempdir().unwrap(); diff --git a/src/main.rs b/src/main.rs index a0cf78b..5c49a25 100644 --- a/src/main.rs +++ b/src/main.rs @@ -675,6 +675,8 @@ fn run() -> Result<()> { .filter(|p| lang.is_empty() || p.has_any_language(&lang)) .collect(); + reject_exporters_from_execution(&target, &affected_projects, &projects)?; + if affected_projects.is_empty() { if output_mode == OutputMode::Json { let output = build_execution_output(&[]); @@ -1254,6 +1256,7 @@ fn run() -> Result<()> { // Select initial projects let initial = select_projects(&run_args, &graph, &projects, &cwd, &workspace_root) .map_err(|e| anyhow::anyhow!("{e}"))?; + reject_exporters_from_execution(&run_args.target, &initial, &projects)?; // Build set of primary projects (originally selected, before expansion) // Only these will run the requested target; dependency projects are included @@ -1507,6 +1510,41 @@ fn handle_init(cwd: &std::path::Path, verbose: bool) -> Result<()> { Ok(()) } +fn reject_exporters_from_execution( + target_name: &str, + primary_projects: &[&DiscoveredProject], + projects: &[DiscoveredProject], +) -> Result<()> { + let project_map = projects + .iter() + .map(|project| (format!("//{}", project.relative_path.display()), project)) + .collect::>(); + let exporter_addresses = projects + .iter() + .flat_map(|project| { + project + .targets + .iter() + .filter(|(_, target)| target.exports_vars()) + .map(move |(name, _)| format!("//{}:{name}", project.relative_path.display())) + }) + .collect::>(); + for project in primary_projects { + let address = format!("//{}:{target_name}", project.relative_path.display()); + let mut closure = HashSet::from([address.clone()]); + collect_target_deps(&address, &project_map, &mut closure); + if let Some(exporter) = closure + .iter() + .find(|address| exporter_addresses.contains(*address)) + { + anyhow::bail!( + "Target '{exporter}' exports runtime variables and can only be run by `aster services up`" + ); + } + } + Ok(()) +} + /// Handle the `watch` command. #[allow(clippy::too_many_arguments)] fn handle_watch( @@ -1579,9 +1617,29 @@ fn handle_watch( return Err(anyhow::anyhow!("{cycle}")); } + let exporter_addresses = projects + .iter() + .flat_map(|project| { + project + .targets + .iter() + .filter(|(_, target)| target.exports_vars()) + .map(move |(name, _)| format!("//{}:{name}", project.relative_path.display())) + }) + .collect::>(); let registry = PluginRegistry::with_all_plugins(); let plan = WatchPlan::build(&resolved, &projects, &graph, ®istry)?; + if let Some(exporter) = plan + .targets + .iter() + .find(|target| exporter_addresses.contains(&target.address)) + { + return Err(anyhow::anyhow!( + "Target '{}' exports runtime variables and can only be run by `aster services up`", + exporter.address + )); + } let workspace_config = WorkspaceConfig::load(workspace_root)?; let ignore = WorkspaceIgnore::build(&workspace_config.watch)?; diff --git a/src/plugins/mod.rs b/src/plugins/mod.rs index 09ecf87..429d5fa 100644 --- a/src/plugins/mod.rs +++ b/src/plugins/mod.rs @@ -43,6 +43,10 @@ pub enum TargetCapability { /// Target can treat warnings as errors /// (e.g., fail build on compiler warnings) WarningsAsErrors, + /// Internal marker for a streaming service target that publishes runtime + /// variables. This is configured by `exports_vars`, not by the public + /// `capabilities` list. + ExportsVars, } /// A build target with its command and dependencies @@ -74,6 +78,13 @@ pub struct Target { pub exclusive_resources: Vec, } +impl Target { + /// Whether this target publishes variable snapshots while supervised. + pub fn exports_vars(&self) -> bool { + self.capabilities.contains(&TargetCapability::ExportsVars) + } +} + /// Context passed to plugins for target detection /// /// Contains all the raw information a plugin needs to determine targets diff --git a/src/targets/resolver.rs b/src/targets/resolver.rs index 9b6b434..7f82086 100644 --- a/src/targets/resolver.rs +++ b/src/targets/resolver.rs @@ -103,11 +103,15 @@ impl TargetResolver { Self::resolve_self_references(&rich.depends_on, project_address) }; - let capabilities = if rich.capabilities.is_empty() { + let mut capabilities = if rich.capabilities.is_empty() { existing.map(|e| e.capabilities.clone()).unwrap_or_default() } else { Self::parse_capabilities(&rich.capabilities) }; + capabilities.remove(&TargetCapability::ExportsVars); + if rich.exports_vars { + capabilities.insert(TargetCapability::ExportsVars); + } let files_glob = rich .files_glob @@ -277,6 +281,7 @@ mod tests { capabilities: capabilities.into_iter().map(|s| s.to_string()).collect(), files_glob: files_glob.map(|s| s.to_string()), stream: false, + exports_vars: false, cache: None, invalidates_cache: false, exclusive_resources: vec![], @@ -517,6 +522,23 @@ mod tests { assert_eq!(integration_target.files_glob, Some("*_test.go".to_string())); } + #[test] + fn test_rich_exporter_marker_survives_alias_resolution() { + let mut custom = HashMap::new(); + let mut exporter = match rich("./publish", vec![], vec![], None) { + TargetConfig::Rich(config) => config, + _ => unreachable!(), + }; + exporter.stream = true; + exporter.exports_vars = true; + custom.insert("exporter".to_string(), TargetConfig::Rich(exporter)); + custom.insert("exporter-alias".to_string(), alias("exporter", vec![])); + + let targets = TargetResolver::resolve(&HashMap::new(), &custom, "//app"); + assert!(targets["exporter"].exports_vars()); + assert!(targets["exporter-alias"].exports_vars()); + } + fn alias(alias_name: &str, depends_on: Vec<&str>) -> TargetConfig { TargetConfig::Alias(crate::config::AliasTargetConfig { alias: alias_name.to_string(), @@ -677,6 +699,7 @@ mod tests { capabilities: vec![], files_glob: None, stream: false, + exports_vars: false, cache: None, invalidates_cache: false, exclusive_resources: vec!["custom_resource".to_string()], diff --git a/tests/dev_services.rs b/tests/dev_services.rs index 7921a99..e0ef6fa 100644 --- a/tests/dev_services.rs +++ b/tests/dev_services.rs @@ -157,6 +157,137 @@ fn terminate_aster(child: &mut std::process::Child) { assert_eq!(status.code(), Some(143)); } +#[test] +fn exported_variables_gate_and_restart_only_direct_consumers() { + let temp = tempfile::tempdir().unwrap(); + let root = temp.path(); + fs::create_dir(root.join(".git")).unwrap(); + for project in ["producer", "consumer", "transitive", "sibling"] { + fs::create_dir(root.join(project)).unwrap(); + fs::write( + root.join(project).join("package.json"), + format!(r#"{{"name":"{project}"}}"#), + ) + .unwrap(); + } + fs::write( + root.join("producer/publish.sh"), + r#"#!/bin/sh +printf '%s\n' '{"TOKEN":"one","PRIVATE":"hidden"}' > "$ASTER_EXPORT_VAR_PATH" +while [ ! -f ../consumer.ready ]; do sleep 0.05; done +sleep 0.2 +printf '%s\n' '{"TOKEN":"two","PRIVATE":"hidden"}' > "$ASTER_EXPORT_VAR_PATH" +sleep 30 +"#, + ) + .unwrap(); + fs::write( + root.join("run.sh"), + r#"#!/bin/sh +name="$1" +printf '%s:%s\n' "$name" "${TOKEN-unset}" >> ../events.log +if [ "$name" = consumer ]; then touch ../consumer.ready; fi +trap 'exit 0' TERM INT +while :; do sleep 1; done +"#, + ) + .unwrap(); + fs::write( + root.join("aster.toml"), + r#" +[dev.services.producer] +target = "//producer:dev" + +[dev.services.consumer] +target = "//consumer:dev" +inherit_env = ["TOKEN"] + +[dev.services.transitive] +target = "//transitive:dev" +inherit_env = ["TOKEN"] + +[dev.services.sibling] +target = "//sibling:dev" +inherit_env = ["TOKEN"] +"#, + ) + .unwrap(); + fs::write( + root.join("producer/aster.toml"), + r#" +[targets.dev] +command = "sh publish.sh" +stream = true +exports_vars = true +"#, + ) + .unwrap(); + fs::write( + root.join("consumer/aster.toml"), + r#" +[targets.dev] +command = "sh ../run.sh consumer" +stream = true +depends_on = ["//producer:dev"] +"#, + ) + .unwrap(); + fs::write( + root.join("transitive/aster.toml"), + r#" +[targets.dev] +command = "sh ../run.sh transitive" +stream = true +depends_on = ["//consumer:dev"] +"#, + ) + .unwrap(); + fs::write( + root.join("sibling/aster.toml"), + r#" +[targets.dev] +command = "sh ../run.sh sibling" +stream = true +"#, + ) + .unwrap(); + + let stdout = fs::File::create(root.join("stdout.log")).unwrap(); + let stderr = fs::File::create(root.join("stderr.log")).unwrap(); + let mut aster = Command::new(env!("CARGO_BIN_EXE_aster")) + .args(["services", "up", "--no-ui", "--no-watch"]) + .env_remove("TOKEN") + .current_dir(root) + .stdout(stdout) + .stderr(stderr) + .spawn() + .unwrap(); + let events = root.join("events.log"); + if !condition_met(Duration::from_secs(15), || { + occurrences(&events, "consumer:one") == 1 + && occurrences(&events, "consumer:two") == 1 + && occurrences(&events, "transitive:unset") == 1 + && occurrences(&events, "sibling:unset") == 1 + }) { + fail_with_process_diagnostics( + &mut aster, + &events, + &root.join("stdout.log"), + &root.join("stderr.log"), + "exported-variable reconciliation did not settle", + ); + } + assert_eq!(occurrences(&events, "transitive:"), 1); + assert_eq!(occurrences(&events, "sibling:"), 1); + assert!(!fs::read_to_string(root.join("stdout.log")) + .unwrap_or_default() + .contains("hidden")); + assert!(!fs::read_to_string(root.join("stderr.log")) + .unwrap_or_default() + .contains("hidden")); + terminate_aster(&mut aster); +} + #[test] fn services_kill_ports_previews_then_clears_configured_listener() { let temp = tempfile::tempdir().unwrap(); From 586d499c266ac1ceffbe33d32131e805baa67db7 Mon Sep 17 00:00:00 2001 From: Arne Roomann-Kurrik Date: Thu, 3 Sep 2026 09:41:44 -0700 Subject: [PATCH 2/6] polish: address round 1 review feedback MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Exporter start failure after ready aborted the supervisor — hold consumers instead, matching process-exit after the first snapshot. - 20s first-snapshot timer included prerequisites and later generations — start the clock when the process is running and abort only before ready. - Colliding export keys unwound the event loop — treat collisions as a hold for the consumer, not a supervisor fatal error. - aster run (and watch/executor spawn) could still execute exporters — reject at the CLI closure and at the spawn seams. - Daemon ready deadline was a single 25s window for serial exporters — refresh it as the supervisor starts each generation. - No tests for the services-up-only gate — added CLI and spawn coverage. Reviewers: grok-native, review-principles --- src/dev/daemon.rs | 79 +++++++++--- src/dev/mod.rs | 5 +- src/dev/process.rs | 3 + src/dev/runner.rs | 276 ++++++++++++++++++++++++++++------------- src/executor/runner.rs | 22 ++++ src/main.rs | 77 ++++++------ src/plugins/mod.rs | 7 ++ src/watch/stream.rs | 18 +++ tests/dev_services.rs | 50 ++++++++ 9 files changed, 397 insertions(+), 140 deletions(-) diff --git a/src/dev/daemon.rs b/src/dev/daemon.rs index 86a6dfd..8bc6ee5 100644 --- a/src/dev/daemon.rs +++ b/src/dev/daemon.rs @@ -3,6 +3,7 @@ use serde::{Deserialize, Serialize}; use std::collections::BTreeMap; use std::path::{Path, PathBuf}; +use std::time::Duration; pub const PROTOCOL_VERSION: u16 = 4; const PREVIOUS_PROTOCOL_VERSION: u16 = 3; @@ -240,8 +241,8 @@ mod platform { const PID_NAME: &str = "daemon.pid"; const LOG_NAME: &str = "daemon.log"; const READY_NAME: &str = "ready.sock"; - // Exporting services own a 20-second first-snapshot deadline. Leave a - // bounded handoff margin for the supervisor to report that result. + // Initial bound until the supervisor can extend the deadline as it starts + // exporters. Each generation refreshes this window; see runner.rs. const START_TIMEOUT: Duration = Duration::from_secs(25); const CONNECT_TIMEOUT: Duration = Duration::from_secs(10); const IDLE_GRACE: Duration = Duration::from_millis(500); @@ -273,6 +274,14 @@ mod platform { services: Vec, ports: BTreeMap, attach_socket: Option, + #[serde(default = "ready_record_default")] + ready: bool, + #[serde(default)] + extend_deadline_secs: Option, + } + + fn ready_record_default() -> bool { + true } pub(super) fn is_serve_invocation() -> bool { @@ -421,14 +430,19 @@ mod platform { Ok(stopped) } - pub(super) fn register_ready( - workspace: &Path, - services: Vec, - ports: BTreeMap, - ) -> DaemonResult<()> { + fn send_ready_record(record: serde_json::Value) -> DaemonResult<()> { let Some(socket) = std::env::var_os(READY_SOCKET_ENV) else { return Ok(()); }; + let datagram = + UnixDatagram::unbound().map_err(internal("failed to create readiness channel"))?; + datagram + .send_to(record.to_string().as_bytes(), socket) + .map_err(internal("failed to register supervisor readiness"))?; + Ok(()) + } + + fn ready_identity(workspace: &Path) -> DaemonResult<(String, String, Option, PathBuf)> { let bundle_id = std::env::var(READY_ID_ENV).map_err(|_| { daemon_error(DaemonErrorCode::InvalidRequest, "missing daemon bundle id") })?; @@ -438,19 +452,38 @@ mod platform { .ok() .filter(|value| !value.is_empty()); let workspace = canonical_directory(workspace, "workspace")?; + Ok((bundle_id, group, display_group, workspace)) + } + + pub(super) fn register_ready( + workspace: &Path, + services: Vec, + ports: BTreeMap, + ) -> DaemonResult<()> { + if std::env::var_os(READY_SOCKET_ENV).is_none() { + return Ok(()); + } + let (bundle_id, group, display_group, workspace) = ready_identity(workspace)?; let attach_socket = std::env::var_os(super::super::session::ATTACH_SOCKET_ENV).map(PathBuf::from); - let record = serde_json::json!({ + send_ready_record(serde_json::json!({ "version": PROTOCOL_VERSION, "bundle_id": bundle_id, "workspace": workspace, "group": group, "display_group": display_group, "supervisor_pid": std::process::id(), "services": services, "ports": ports, "attach_socket": attach_socket, - }); - let datagram = - UnixDatagram::unbound().map_err(internal("failed to create readiness channel"))?; - datagram - .send_to(record.to_string().as_bytes(), socket) - .map_err(internal("failed to register supervisor readiness"))?; - Ok(()) + })) + } + + pub(super) fn extend_ready_deadline(workspace: &Path, extra: Duration) -> DaemonResult<()> { + if std::env::var_os(READY_SOCKET_ENV).is_none() { + return Ok(()); + } + let (bundle_id, group, display_group, workspace) = ready_identity(workspace)?; + send_ready_record(serde_json::json!({ + "version": PROTOCOL_VERSION, "bundle_id": bundle_id, "workspace": workspace, + "group": group, "display_group": display_group, "supervisor_pid": std::process::id(), + "services": [], "ports": {}, "attach_socket": Option::::None, + "ready": false, "extend_deadline_secs": extra.as_secs(), + })) } pub(super) struct RuntimePaths { @@ -1099,6 +1132,14 @@ mod platform { { continue; } + if !record.ready { + if let Some(seconds) = record.extend_deadline_secs { + if bundle.descriptor.state == BundleState::Starting { + bundle.ready_deadline = Instant::now() + Duration::from_secs(seconds); + } + } + continue; + } bundle.descriptor.state = BundleState::Running; bundle.descriptor.services = record.services; bundle.descriptor.ports = record.ports; @@ -1338,6 +1379,7 @@ mod platform { #[cfg(not(unix))] mod platform { use super::*; + use std::time::Duration; fn unsupported() -> DaemonResult { Err(daemon_error( DaemonErrorCode::UnsupportedPlatform, @@ -1375,6 +1417,9 @@ mod platform { ) -> DaemonResult<()> { Ok(()) } + pub(super) fn extend_ready_deadline(_: &Path, _: Duration) -> DaemonResult<()> { + Ok(()) + } } pub fn ping_daemon() -> DaemonResult { @@ -1416,6 +1461,10 @@ pub fn register_supervisor_ready( platform::register_ready(workspace, services, ports) } +pub fn register_supervisor_extend_deadline(workspace: &Path, extra: Duration) -> DaemonResult<()> { + platform::extend_ready_deadline(workspace, extra) +} + fn protocol_mismatch() -> DaemonError { daemon_error( DaemonErrorCode::UnsupportedProtocol, diff --git a/src/dev/mod.rs b/src/dev/mod.rs index 5d11013..25c4e98 100644 --- a/src/dev/mod.rs +++ b/src/dev/mod.rs @@ -21,7 +21,10 @@ pub use daemon::{ SUPERVISOR_ENV, }; #[doc(hidden)] -pub use daemon::{is_internal_serve_invocation, register_supervisor_ready, serve_from_environment}; +pub use daemon::{ + is_internal_serve_invocation, register_supervisor_extend_deadline, register_supervisor_ready, + serve_from_environment, +}; pub use log_files::show_service_logs; pub use plan::{ resolve_dev_plan, resolve_dev_ports, resolve_static_dev_ports, DevPlan, ServicePlan, diff --git a/src/dev/process.rs b/src/dev/process.rs index 4b760f4..da9fd4a 100644 --- a/src/dev/process.rs +++ b/src/dev/process.rs @@ -66,6 +66,9 @@ impl ServiceProcess { env: &HashMap, log_senders: ProcessLogSenders, ) -> Result { + if target.exports_vars() && !env.contains_key(super::export_vars::EXPORT_PATH_ENV) { + anyhow::bail!(crate::plugins::exporter_requires_services_up(service)); + } let parsed = parse_command(&target.command)?; let working_dir = target.working_dir.as_deref().unwrap_or(project_root); let mut command = Command::new(&parsed.program); diff --git a/src/dev/runner.rs b/src/dev/runner.rs index b3c09ea..4008d38 100644 --- a/src/dev/runner.rs +++ b/src/dev/runner.rs @@ -61,6 +61,20 @@ struct ResolvedExports { sources: HashSet, } +#[derive(Debug)] +enum ExportResolution { + Ready(ResolvedExports), + Pending, + Conflict(String), +} + +const FIRST_SNAPSHOT_TIMEOUT: Duration = Duration::from_secs(20); +const READY_HANDOFF_SECS: u64 = 5; +const EXPORTER_PREREQ_BUDGET_SECS: u64 = 60; +const READY_PROGRESS_SECS: u64 = + FIRST_SNAPSHOT_TIMEOUT.as_secs() + READY_HANDOFF_SECS + EXPORTER_PREREQ_BUDGET_SECS; +const SNAPSHOT_PROGRESS_SECS: u64 = FIRST_SNAPSHOT_TIMEOUT.as_secs() + READY_HANDOFF_SECS; + enum StartOutcome { Running { process: ServiceProcess, @@ -223,13 +237,20 @@ pub fn run_dev( .iter() .find(|service| service.name == result.service) .expect("start result service exists"); - let current = resolve_exported_environment(service, &producers)?; - if current.as_ref().map(|resolved| &resolved.values) - != Some(&result.resolved_exports.values) - { + let current = resolve_exported_environment(service, &producers); + let keep_running = match ¤t { + ExportResolution::Ready(resolved) => { + resolved.values == result.resolved_exports.values + } + ExportResolution::Pending | ExportResolution::Conflict(_) => false, + }; + if !keep_running { process.terminate(Duration::from_secs(3)); #[cfg(unix)] drop(export_endpoint); + if let ExportResolution::Conflict(message) = current { + emit_system(&system_tx, &result.service, message, true); + } queue_restart( &mut pending_starts, &result.service, @@ -238,6 +259,9 @@ pub fn run_dev( &mut dashboard, ); } else { + let ExportResolution::Ready(current) = current else { + unreachable!("keep_running requires a ready export resolution"); + }; let runtime = runtimes .get_mut(&result.service) .expect("start result service exists"); @@ -248,9 +272,18 @@ pub fn run_dev( } runtime.applied_exports = result.resolved_exports.values; runtime.applied_sources = result.resolved_exports.sources; - runtime.applied_epochs = current - .map(|resolved| resolved.epochs) - .unwrap_or(result.resolved_exports.epochs); + runtime.applied_epochs = current.epochs; + if let Some(state) = producers.get_mut(&result.service) { + if state.values.is_none() { + state.waiting_since = Some(Instant::now()); + } + } + if !supervisor_registered && producers.contains_key(&result.service) { + extend_ready_deadline( + workspace_root, + Duration::from_secs(SNAPSHOT_PROGRESS_SECS), + ); + } set_service_state( &mut dashboard, #[cfg(unix)] @@ -261,13 +294,22 @@ pub fn run_dev( } } StartOutcome::Stopped => { - if producers - .get(&result.service) - .is_some_and(|producer| producer.values.is_none()) - { - anyhow::bail!( - "variable exporter '{}' exited before publishing its initial snapshot", - result.service + if producers.contains_key(&result.service) { + if !supervisor_registered { + anyhow::bail!( + "variable exporter '{}' exited before publishing its initial snapshot", + result.service + ); + } + stop_dependent_services( + &result.service, + &plan, + &mut runtimes, + &mut pending_starts, + &system_tx, + &mut dashboard, + #[cfg(unix)] + session.as_ref(), ); } set_service_state( @@ -310,28 +352,48 @@ pub fn run_dev( false, ); for service in &plan.services { - if service.name == producer - || !service.dependencies.contains(&producer) - || runtimes[&service.name].process.is_none() + if service.name == producer || !service.dependencies.contains(&producer) { continue; } - let desired = resolve_exported_environment(service, &producers)?; - if desired.as_ref().map(|resolved| &resolved.values) - != Some(&runtimes[&service.name].applied_exports) - { - queue_restart( - &mut pending_starts, - &service.name, - "effective exported environment changed", - &system_tx, - &mut dashboard, - ); - } else if let Some(resolved) = desired { - let runtime = - runtimes.get_mut(&service.name).expect("runtime exists"); - runtime.applied_epochs = resolved.epochs; - runtime.applied_sources = resolved.sources; + match resolve_exported_environment(service, &producers) { + ExportResolution::Ready(resolved) + if runtimes[&service.name].process.is_some() + && resolved.values + == runtimes[&service.name].applied_exports => + { + let runtime = + runtimes.get_mut(&service.name).expect("runtime exists"); + runtime.applied_epochs = resolved.epochs; + runtime.applied_sources = resolved.sources; + } + ExportResolution::Ready(_) + if runtimes[&service.name].process.is_some() => + { + queue_restart( + &mut pending_starts, + &service.name, + "effective exported environment changed", + &system_tx, + &mut dashboard, + ); + } + ExportResolution::Ready(_) | ExportResolution::Pending => {} + ExportResolution::Conflict(message) => { + emit_system(&system_tx, &service.name, message, true); + if runtimes[&service.name].process.is_some() { + hold_service( + &service.name, + "exported variable sources conflict", + &mut runtimes, + &mut pending_starts, + &system_tx, + &mut dashboard, + #[cfg(unix)] + session.as_ref(), + ); + } + } } } } @@ -378,14 +440,16 @@ pub fn run_dev( supervisor_registered = true; } - if let Some((name, _)) = producers.iter().find(|(_, producer)| { - producer.waiting_since.is_some_and(|started| { - started.elapsed() >= Duration::from_secs(20) && producer.values.is_none() - }) - }) { - anyhow::bail!( - "variable exporter '{name}' did not publish an initial snapshot within 20 seconds" - ); + if !supervisor_registered { + if let Some((name, _)) = producers.iter().find(|(_, producer)| { + producer.waiting_since.is_some_and(|started| { + started.elapsed() >= FIRST_SNAPSHOT_TIMEOUT && producer.values.is_none() + }) + }) { + anyhow::bail!( + "variable exporter '{name}' did not publish an initial snapshot within 20 seconds" + ); + } } if active_start.is_none() { let mut startable = None; @@ -395,7 +459,7 @@ pub fn run_dev( .iter() .find(|service| service.name == *name) .expect("queued service exists"); - if service_is_startable(service, &runtimes, &producers)? { + if service_is_startable(service, &runtimes, &producers) { startable = Some(index); break; } @@ -428,8 +492,24 @@ pub fn run_dev( let before_start = options .watch .then(|| suppressed_path_snapshot(&plan, workspace_root, &ignore)); - let resolved_exports = - resolve_exported_environment(service, &producers)?.unwrap_or_default(); + let resolved_exports = match resolve_exported_environment(service, &producers) { + ExportResolution::Ready(resolved) => resolved, + ExportResolution::Conflict(message) => { + emit_system(&system_tx, &name, message, true); + pending_starts.push_front((name, reason)); + continue; + } + ExportResolution::Pending => { + pending_starts.push_front((name, reason)); + continue; + } + }; + if !supervisor_registered { + extend_ready_deadline( + workspace_root, + Duration::from_secs(READY_PROGRESS_SECS), + ); + } let handle = begin_start_service( service, workspace_root, @@ -829,10 +909,14 @@ fn register_ready(workspace_root: &Path, plan: &DevPlan) -> Result<()> { Ok(()) } +fn extend_ready_deadline(workspace_root: &Path, extra: Duration) { + let _ = super::daemon::register_supervisor_extend_deadline(workspace_root, extra); +} + fn resolve_exported_environment( service: &ServicePlan, producers: &HashMap, -) -> Result> { +) -> ExportResolution { let mut environment = HashMap::new(); let mut epochs = HashMap::new(); let mut sources = HashMap::::new(); @@ -842,10 +926,10 @@ fn resolve_exported_environment( continue; }; if !producer.healthy { - return Ok(None); + return ExportResolution::Pending; } let Some(values) = producer.values.as_ref() else { - return Ok(None); + return ExportResolution::Pending; }; for name in &service.inherit_env { if service.protected_env.contains(name) { @@ -855,10 +939,10 @@ fn resolve_exported_environment( continue; }; if let Some(existing) = sources.insert(name.clone(), dependency.clone()) { - anyhow::bail!( + return ExportResolution::Conflict(format!( "service '{}' consumes exported variable '{name}' from both '{existing}' and '{dependency}'", service.name - ); + )); } environment.insert(name.clone(), value.clone()); used_sources.insert(dependency.clone()); @@ -866,26 +950,58 @@ fn resolve_exported_environment( epochs.insert(dependency.clone(), producer.epoch); } - Ok(Some(ResolvedExports { + ExportResolution::Ready(ResolvedExports { values: environment, epochs, sources: used_sources, - })) + }) } fn service_is_startable( service: &ServicePlan, runtimes: &HashMap, producers: &HashMap, -) -> Result { +) -> bool { if service.dependencies.iter().any(|dependency| { runtimes .get(dependency) .is_none_or(|runtime| runtime.process.is_none()) }) { - return Ok(false); + return false; } - Ok(resolve_exported_environment(service, producers)?.is_some()) + matches!( + resolve_exported_environment(service, producers), + ExportResolution::Ready(_) + ) +} + +#[allow(clippy::too_many_arguments)] +fn hold_service( + name: &str, + reason: &str, + runtimes: &mut HashMap, + pending: &mut VecDeque<(String, String)>, + system_tx: &std::sync::mpsc::Sender, + dashboard: &mut Dashboard, + #[cfg(unix)] session: Option<&SessionServer>, +) { + let runtime = runtimes.get_mut(name).expect("runtime exists"); + #[cfg(unix)] + runtime.export_endpoint.take(); + if let Some(mut process) = runtime.process.take() { + process.terminate(Duration::from_secs(3)); + } + runtime.applied_exports.clear(); + runtime.applied_epochs.clear(); + runtime.applied_sources.clear(); + set_service_state( + dashboard, + #[cfg(unix)] + session, + name, + ServiceState::Restarting, + ); + queue_restart(pending, name, reason, system_tx, dashboard); } #[allow(clippy::too_many_arguments)] @@ -918,34 +1034,17 @@ fn stop_dependent_services( if !affected.contains(&service.name) { continue; } - let runtime = runtimes.get_mut(&service.name).expect("runtime exists"); - #[cfg(unix)] - runtime.export_endpoint.take(); - if let Some(mut process) = runtime.process.take() { - process.terminate(Duration::from_secs(3)); - } - runtime.applied_exports.clear(); - runtime.applied_epochs.clear(); - runtime.applied_sources.clear(); - set_service_state( + hold_service( + &service.name, + "required exported variables unavailable", + runtimes, + pending, + system_tx, dashboard, #[cfg(unix)] session, - &service.name, - ServiceState::Restarting, ); } - for service in &plan.services { - if affected.contains(&service.name) { - queue_restart( - pending, - &service.name, - "required exported variables unavailable", - system_tx, - dashboard, - ); - } - } } fn set_service_state( @@ -1084,7 +1183,7 @@ fn begin_start_service( producer.epoch = producer.epoch.saturating_add(1); producer.values = None; producer.healthy = false; - producer.waiting_since = Some(Instant::now()); + producer.waiting_since = None; } let state = if reason == "initial start" { ServiceState::Starting @@ -1542,9 +1641,10 @@ mod tests { ready_producer(4, &[("TOKEN", "masked"), ("PRIVATE", "selected")]), )]); - let resolved = resolve_exported_environment(&service, &producers) - .unwrap() - .unwrap(); + let ExportResolution::Ready(resolved) = resolve_exported_environment(&service, &producers) + else { + panic!("expected a ready exported environment"); + }; assert_eq!( resolved.values, HashMap::from([("PRIVATE".to_string(), "selected".to_string())]) @@ -1560,15 +1660,17 @@ mod tests { ("first".to_string(), ready_producer(1, &[("TOKEN", "one")])), ("second".to_string(), ready_producer(1, &[("TOKEN", "two")])), ]); - let error = resolve_exported_environment(&service, &producers) - .unwrap_err() - .to_string(); + let ExportResolution::Conflict(error) = resolve_exported_environment(&service, &producers) + else { + panic!("expected a colliding exported environment"); + }; assert!(error.contains("from both 'first' and 'second'")); producers.get_mut("second").unwrap().healthy = false; - assert!(resolve_exported_environment(&service, &producers) - .unwrap() - .is_none()); + assert!(matches!( + resolve_exported_environment(&service, &producers), + ExportResolution::Pending + )); } #[test] diff --git a/src/executor/runner.rs b/src/executor/runner.rs index 0b79b37..96724c4 100644 --- a/src/executor/runner.rs +++ b/src/executor/runner.rs @@ -182,6 +182,9 @@ impl<'a> Executor<'a> { return Err(format!("No '{target}' target defined for {project_addr}")); } }; + if target_def.exports_vars() { + return Err(crate::plugins::exporter_requires_services_up(&target_addr)); + } let command = &target_def.command; // Use target's working_dir if set, otherwise use project root @@ -328,6 +331,25 @@ impl<'a> Executor<'a> { return Vec::new(); } + if let Some(address) = targets_to_run.iter().find(|address| { + parse_target_address(address) + .and_then(|(project_addr, target_name)| { + project_map + .get(&project_addr) + .and_then(|project| project.targets.get(&target_name)) + }) + .is_some_and(|target| target.exports_vars()) + }) { + return vec![ExecutionResult { + address: address.clone(), + success: false, + skipped: false, + cached: false, + output: crate::plugins::exporter_requires_services_up(address), + duration_ms: 0, + }]; + } + if self.output_mode != OutputMode::Json { let log_store = LogStore::new(self.workspace_root); if let Err(error) = log_store.clear_latest() { diff --git a/src/main.rs b/src/main.rs index 5c49a25..79a6232 100644 --- a/src/main.rs +++ b/src/main.rs @@ -27,7 +27,7 @@ use aster::git::{ AffectedIgnore, }; use aster::graph::{build_graph, build_target_graph, find_cycle, format_path}; -use aster::plugins::{PluginRegistry, Target, TargetCapability}; +use aster::plugins::{exporter_requires_services_up, PluginRegistry, Target, TargetCapability}; use chrono::{DateTime, Utc}; use globset::{Glob, GlobMatcher}; use std::collections::{HashMap, HashSet}; @@ -1038,6 +1038,13 @@ fn run() -> Result<()> { let primary_targets: std::collections::HashSet = filtered_targets.into_iter().collect(); + let mut closure = primary_targets.clone(); + if !no_deps { + for address in &primary_targets { + collect_target_deps(address, &project_map, &mut closure); + } + } + reject_exporter_addresses(&closure, &projects)?; let all_project_refs: Vec<_> = projects.iter().collect(); let executor = Executor::with_all_options(&workspace_root, output_mode, full_logs, !cli.no_cache); @@ -1510,6 +1517,33 @@ fn handle_init(cwd: &std::path::Path, verbose: bool) -> Result<()> { Ok(()) } +fn exporter_addresses(projects: &[DiscoveredProject]) -> HashSet { + projects + .iter() + .flat_map(|project| { + project + .targets + .iter() + .filter(|(_, target)| target.exports_vars()) + .map(move |(name, _)| format!("//{}:{name}", project.relative_path.display())) + }) + .collect() +} + +fn reject_exporter_addresses<'a, I>(addresses: I, projects: &[DiscoveredProject]) -> Result<()> +where + I: IntoIterator, +{ + let exporters = exporter_addresses(projects); + if let Some(exporter) = addresses + .into_iter() + .find(|address| exporters.contains(*address)) + { + anyhow::bail!("{}", exporter_requires_services_up(exporter)); + } + Ok(()) +} + fn reject_exporters_from_execution( target_name: &str, primary_projects: &[&DiscoveredProject], @@ -1519,28 +1553,11 @@ fn reject_exporters_from_execution( .iter() .map(|project| (format!("//{}", project.relative_path.display()), project)) .collect::>(); - let exporter_addresses = projects - .iter() - .flat_map(|project| { - project - .targets - .iter() - .filter(|(_, target)| target.exports_vars()) - .map(move |(name, _)| format!("//{}:{name}", project.relative_path.display())) - }) - .collect::>(); for project in primary_projects { let address = format!("//{}:{target_name}", project.relative_path.display()); let mut closure = HashSet::from([address.clone()]); collect_target_deps(&address, &project_map, &mut closure); - if let Some(exporter) = closure - .iter() - .find(|address| exporter_addresses.contains(*address)) - { - anyhow::bail!( - "Target '{exporter}' exports runtime variables and can only be run by `aster services up`" - ); - } + reject_exporter_addresses(&closure, projects)?; } Ok(()) } @@ -1617,29 +1634,15 @@ fn handle_watch( return Err(anyhow::anyhow!("{cycle}")); } - let exporter_addresses = projects - .iter() - .flat_map(|project| { - project - .targets - .iter() - .filter(|(_, target)| target.exports_vars()) - .map(move |(name, _)| format!("//{}:{name}", project.relative_path.display())) - }) - .collect::>(); let registry = PluginRegistry::with_all_plugins(); let plan = WatchPlan::build(&resolved, &projects, &graph, ®istry)?; - if let Some(exporter) = plan + let planned = plan .targets .iter() - .find(|target| exporter_addresses.contains(&target.address)) - { - return Err(anyhow::anyhow!( - "Target '{}' exports runtime variables and can only be run by `aster services up`", - exporter.address - )); - } + .map(|target| target.address.clone()) + .collect::>(); + reject_exporter_addresses(&planned, &projects)?; let workspace_config = WorkspaceConfig::load(workspace_root)?; let ignore = WorkspaceIgnore::build(&workspace_config.watch)?; diff --git a/src/plugins/mod.rs b/src/plugins/mod.rs index 429d5fa..5bab54e 100644 --- a/src/plugins/mod.rs +++ b/src/plugins/mod.rs @@ -85,6 +85,13 @@ impl Target { } } +/// Error text used when an exporter target is invoked outside `aster services up`. +pub fn exporter_requires_services_up(address: &str) -> String { + format!( + "Target '{address}' exports runtime variables and can only be run by `aster services up`" + ) +} + /// Context passed to plugins for target detection /// /// Contains all the raw information a plugin needs to determine targets diff --git a/src/watch/stream.rs b/src/watch/stream.rs index ac7842a..43de438 100644 --- a/src/watch/stream.rs +++ b/src/watch/stream.rs @@ -106,6 +106,9 @@ impl StreamSupervisor { if self.children.contains_key(target_addr) { return Err(anyhow!("{target_addr} already running")); } + if target.exports_vars() { + anyhow::bail!(crate::plugins::exporter_requires_services_up(target_addr)); + } let child = spawn_child(target, project_root) .with_context(|| format!("failed to spawn {target_addr}"))?; @@ -218,6 +221,7 @@ fn spawn_child(target: &Target, project_root: &Path) -> Result { #[cfg(test)] mod tests { use super::*; + use crate::plugins::TargetCapability; use std::path::PathBuf; fn mk_target(cmd: &str) -> Target { @@ -231,6 +235,20 @@ mod tests { tempfile::tempdir().unwrap() } + #[test] + fn spawn_rejects_variable_exporters() { + let dir = tempdir(); + let mut sup = StreamSupervisor::new(Duration::from_secs(1)); + let mut target = mk_target("sleep 5"); + target.capabilities.insert(TargetCapability::ExportsVars); + let error = sup + .spawn("//a:dev", &target, dir.path()) + .unwrap_err() + .to_string(); + assert!(error.contains("aster services up")); + assert!(!sup.is_running("//a:dev")); + } + #[test] fn spawn_and_poll_running_child() { let dir = tempdir(); diff --git a/tests/dev_services.rs b/tests/dev_services.rs index e0ef6fa..ed47a3c 100644 --- a/tests/dev_services.rs +++ b/tests/dev_services.rs @@ -288,6 +288,56 @@ stream = true terminate_aster(&mut aster); } +#[test] +fn exported_variable_targets_are_rejected_outside_services_up() { + let temp = tempfile::tempdir().unwrap(); + let root = temp.path(); + fs::create_dir(root.join(".git")).unwrap(); + fs::create_dir(root.join("producer")).unwrap(); + fs::write(root.join("producer/package.json"), r#"{"name":"producer"}"#).unwrap(); + fs::write( + root.join("producer/aster.toml"), + r#" +[targets.dev] +command = "sleep 30" +stream = true +exports_vars = true +"#, + ) + .unwrap(); + fs::write( + root.join("aster.toml"), + r#" +[dev.services.producer] +target = "//producer:dev" +"#, + ) + .unwrap(); + + let cases: &[&[&str]] = &[ + &["run", "//producer:dev"], + &["dev", "//producer"], + &["watch", "//producer:dev", "--no-initial"], + ]; + for args in cases { + let output = Command::new(env!("CARGO_BIN_EXE_aster")) + .args(*args) + .current_dir(root) + .output() + .unwrap(); + let stdout = String::from_utf8_lossy(&output.stdout); + let stderr = String::from_utf8_lossy(&output.stderr); + assert!( + !output.status.success(), + "{args:?} should reject exporter targets\nstdout={stdout}\nstderr={stderr}" + ); + assert!( + stdout.contains("aster services up") || stderr.contains("aster services up"), + "{args:?} should name the services-up restriction\nstdout={stdout}\nstderr={stderr}" + ); + } +} + #[test] fn services_kill_ports_previews_then_clears_configured_listener() { let temp = tempfile::tempdir().unwrap(); From 54858dbac5338fca1213b12210b8562c16a23db4 Mon Sep 17 00:00:00 2001 From: Arne Roomann-Kurrik Date: Thu, 3 Sep 2026 10:01:18 -0700 Subject: [PATCH 3/6] polish: address round 2 review feedback MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Ready-deadline refresh used protocol 4 records an old daemon treats as ready — bump to protocol 5 so mixed binaries replace the daemon. - Invalid snapshot reasons forwarded serde text that can include payload values — keep parser errors generic so secrets stay out of logs. - Invalid data after a first snapshot left ready hung with no 20s abort — re-arm the wait and time out on unhealthy, not only missing values. - Daemon deadline was not refreshed while a start generation ran — keep extending it while a start is in flight. - Non-Unix rejected any workspace exporter, not just selected services — gate Unix support on the selected plan. - --no-deps still walked exporter dependencies — skip that closure when the exporter would not run. Reviewers: grok-native, codex-cli, review-principles --- src/dev/daemon.rs | 4 ++-- src/dev/export_vars.rs | 22 ++++++++++++++++++++-- src/dev/plan.rs | 2 +- src/dev/runner.rs | 13 ++++++++++++- src/main.rs | 14 +++++++++++--- tests/dev_services.rs | 26 ++++++++++++++++++++++++++ 6 files changed, 72 insertions(+), 9 deletions(-) diff --git a/src/dev/daemon.rs b/src/dev/daemon.rs index 8bc6ee5..e6f83f1 100644 --- a/src/dev/daemon.rs +++ b/src/dev/daemon.rs @@ -5,8 +5,8 @@ use std::collections::BTreeMap; use std::path::{Path, PathBuf}; use std::time::Duration; -pub const PROTOCOL_VERSION: u16 = 4; -const PREVIOUS_PROTOCOL_VERSION: u16 = 3; +pub const PROTOCOL_VERSION: u16 = 5; +const PREVIOUS_PROTOCOL_VERSION: u16 = 4; pub const DEFAULT_GROUP: &str = "__aster_default__"; const SERVE_ENV: &str = "ASTER_INTERNAL_DAEMON_SERVE"; const READY_SOCKET_ENV: &str = "ASTER_INTERNAL_DAEMON_READY_SOCKET"; diff --git a/src/dev/export_vars.rs b/src/dev/export_vars.rs index 02419c4..bbc30c2 100644 --- a/src/dev/export_vars.rs +++ b/src/dev/export_vars.rs @@ -86,10 +86,10 @@ fn parse_snapshot(bytes: &[u8]) -> Result, String> { let mut deserializer = serde_json::Deserializer::from_slice(bytes); let values = serde::de::Deserializer::deserialize_map(&mut deserializer, SnapshotVisitor) - .map_err(|error| format!("invalid JSON snapshot: {error}"))?; + .map_err(|_| "invalid JSON snapshot".to_string())?; deserializer .end() - .map_err(|error| format!("trailing JSON data: {error}"))?; + .map_err(|_| "trailing JSON data".to_string())?; for (name, value) in &values { if !valid_environment_name(name) { return Err(format!("invalid environment variable name '{name}'")); @@ -335,6 +335,24 @@ mod tests { } } + #[test] + fn parser_errors_do_not_include_payload_values() { + for input in [ + b"\"leaked-secret-token\"\n".as_slice(), + b"{\"TOKEN\":\"leaked-secret-token\"} extra\n".as_slice(), + ] { + let error = JsonlParser::default() + .push(input) + .pop() + .unwrap() + .unwrap_err(); + assert!( + !error.contains("leaked-secret-token"), + "parser error leaked payload: {error}" + ); + } + } + #[cfg(unix)] #[test] fn endpoint_is_private_and_publishes_without_logging_values() { diff --git a/src/dev/plan.rs b/src/dev/plan.rs index b36864b..9df5f64 100644 --- a/src/dev/plan.rs +++ b/src/dev/plan.rs @@ -351,7 +351,7 @@ fn order_service_graph( }) .collect::>(); #[cfg(not(unix))] - if !exporter_addresses.is_empty() { + if services.iter().any(|service| service.target.exports_vars()) { bail!("exports_vars targets require the Unix service supervisor"); } let mut instances: HashMap> = HashMap::new(); diff --git a/src/dev/runner.rs b/src/dev/runner.rs index 4008d38..44aafa5 100644 --- a/src/dev/runner.rs +++ b/src/dev/runner.rs @@ -191,6 +191,7 @@ pub fn run_dev( let watch_debounce = Duration::from_millis(config.watch.debounce_ms.unwrap_or(300)); let mut pending_watch_paths = Vec::new(); let mut watch_deadline: Option = None; + let mut last_ready_extend = Instant::now(); let mut quitting = false; if producers.is_empty() { @@ -283,6 +284,7 @@ pub fn run_dev( workspace_root, Duration::from_secs(SNAPSHOT_PROGRESS_SECS), ); + last_ready_extend = Instant::now(); } set_service_state( &mut dashboard, @@ -410,6 +412,9 @@ pub fn run_dev( } state.epoch = state.epoch.saturating_add(1); state.healthy = false; + if state.waiting_since.is_none() { + state.waiting_since = Some(Instant::now()); + } emit_system( &system_tx, &producer, @@ -441,9 +446,14 @@ pub fn run_dev( } if !supervisor_registered { + if active_start.is_some() && last_ready_extend.elapsed() >= Duration::from_secs(10) + { + extend_ready_deadline(workspace_root, Duration::from_secs(READY_PROGRESS_SECS)); + last_ready_extend = Instant::now(); + } if let Some((name, _)) = producers.iter().find(|(_, producer)| { producer.waiting_since.is_some_and(|started| { - started.elapsed() >= FIRST_SNAPSHOT_TIMEOUT && producer.values.is_none() + started.elapsed() >= FIRST_SNAPSHOT_TIMEOUT && !producer.healthy }) }) { anyhow::bail!( @@ -509,6 +519,7 @@ pub fn run_dev( workspace_root, Duration::from_secs(READY_PROGRESS_SECS), ); + last_ready_extend = Instant::now(); } let handle = begin_start_service( service, diff --git a/src/main.rs b/src/main.rs index 79a6232..120cc67 100644 --- a/src/main.rs +++ b/src/main.rs @@ -675,7 +675,7 @@ fn run() -> Result<()> { .filter(|p| lang.is_empty() || p.has_any_language(&lang)) .collect(); - reject_exporters_from_execution(&target, &affected_projects, &projects)?; + reject_exporters_from_execution(&target, &affected_projects, &projects, true)?; if affected_projects.is_empty() { if output_mode == OutputMode::Json { @@ -1263,7 +1263,12 @@ fn run() -> Result<()> { // Select initial projects let initial = select_projects(&run_args, &graph, &projects, &cwd, &workspace_root) .map_err(|e| anyhow::anyhow!("{e}"))?; - reject_exporters_from_execution(&run_args.target, &initial, &projects)?; + reject_exporters_from_execution( + &run_args.target, + &initial, + &projects, + !run_args.no_deps, + )?; // Build set of primary projects (originally selected, before expansion) // Only these will run the requested target; dependency projects are included @@ -1548,6 +1553,7 @@ fn reject_exporters_from_execution( target_name: &str, primary_projects: &[&DiscoveredProject], projects: &[DiscoveredProject], + expand_deps: bool, ) -> Result<()> { let project_map = projects .iter() @@ -1556,7 +1562,9 @@ fn reject_exporters_from_execution( for project in primary_projects { let address = format!("//{}:{target_name}", project.relative_path.display()); let mut closure = HashSet::from([address.clone()]); - collect_target_deps(&address, &project_map, &mut closure); + if expand_deps { + collect_target_deps(&address, &project_map, &mut closure); + } reject_exporter_addresses(&closure, projects)?; } Ok(()) diff --git a/tests/dev_services.rs b/tests/dev_services.rs index ed47a3c..fd9de4b 100644 --- a/tests/dev_services.rs +++ b/tests/dev_services.rs @@ -302,6 +302,18 @@ fn exported_variable_targets_are_rejected_outside_services_up() { command = "sleep 30" stream = true exports_vars = true +"#, + ) + .unwrap(); + fs::create_dir(root.join("consumer")).unwrap(); + fs::write(root.join("consumer/package.json"), r#"{"name":"consumer"}"#).unwrap(); + fs::write( + root.join("consumer/aster.toml"), + r#" +[targets.dev] +command = "echo ok" +stream = true +depends_on = ["//producer:dev"] "#, ) .unwrap(); @@ -318,6 +330,8 @@ target = "//producer:dev" &["run", "//producer:dev"], &["dev", "//producer"], &["watch", "//producer:dev", "--no-initial"], + &["run", "//consumer:dev"], + &["dev", "//consumer"], ]; for args in cases { let output = Command::new(env!("CARGO_BIN_EXE_aster")) @@ -336,6 +350,18 @@ target = "//producer:dev" "{args:?} should name the services-up restriction\nstdout={stdout}\nstderr={stderr}" ); } + + let output = Command::new(env!("CARGO_BIN_EXE_aster")) + .args(["run", "//consumer:dev", "--no-deps"]) + .current_dir(root) + .output() + .unwrap(); + assert!( + output.status.success(), + "--no-deps should still run a non-exporter primary\nstdout={}\nstderr={}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); } #[test] From 2f36b4e7b24b1fc090d2f199cb54a5a863a7aa22 Mon Sep 17 00:00:00 2001 From: Arne Roomann-Kurrik Date: Thu, 3 Sep 2026 10:23:55 -0700 Subject: [PATCH 4/6] polish: address round 3 review feedback MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Homogeneous --no-deps still collected same-project exporter deps and failed the run — skip dependency expansion in execute_internal too. - Pre-ready abort still said "initial snapshot" after invalid-after- snapshot — report that the exporter did not become healthy. Reviewers: grok-native, review-principles --- src/dev/runner.rs | 2 +- src/executor/runner.rs | 18 ++++++++++++++---- src/main.rs | 11 ++++++++++- tests/dev_services.rs | 33 ++++++++++++++++++++++----------- 4 files changed, 47 insertions(+), 17 deletions(-) diff --git a/src/dev/runner.rs b/src/dev/runner.rs index 44aafa5..fe889d4 100644 --- a/src/dev/runner.rs +++ b/src/dev/runner.rs @@ -457,7 +457,7 @@ pub fn run_dev( }) }) { anyhow::bail!( - "variable exporter '{name}' did not publish an initial snapshot within 20 seconds" + "variable exporter '{name}' did not become healthy within 20 seconds" ); } } diff --git a/src/executor/runner.rs b/src/executor/runner.rs index 96724c4..460124c 100644 --- a/src/executor/runner.rs +++ b/src/executor/runner.rs @@ -144,8 +144,9 @@ impl<'a> Executor<'a> { projects: &[&DiscoveredProject], _graph: &ProjectGraph, primary_projects: Option<&HashSet>, + expand_deps: bool, ) -> Vec { - self.execute_internal(target, projects, None, primary_projects) + self.execute_internal(target, projects, None, primary_projects, expand_deps) } /// Execute a target with command overrides for specific targets @@ -158,8 +159,15 @@ impl<'a> Executor<'a> { projects: &[&DiscoveredProject], command_overrides: &HashMap, primary_projects: Option<&HashSet>, + expand_deps: bool, ) -> Vec { - self.execute_internal(target, projects, Some(command_overrides), primary_projects) + self.execute_internal( + target, + projects, + Some(command_overrides), + primary_projects, + expand_deps, + ) } /// Execute a target with streaming output (for long-running processes like dev servers) @@ -280,6 +288,7 @@ impl<'a> Executor<'a> { projects: &[&DiscoveredProject], command_overrides: Option<&HashMap>, primary_projects: Option<&HashSet>, + expand_deps: bool, ) -> Vec { if projects.is_empty() { return Vec::new(); @@ -310,8 +319,9 @@ impl<'a> Executor<'a> { // Add the requested target targets_to_run.insert(target_addr.clone()); - // Recursively collect target dependencies - collect_target_deps(&target_addr, &project_map, &mut targets_to_run); + if expand_deps { + collect_target_deps(&target_addr, &project_map, &mut targets_to_run); + } } } diff --git a/src/main.rs b/src/main.rs index 120cc67..a6400ff 100644 --- a/src/main.rs +++ b/src/main.rs @@ -952,9 +952,16 @@ fn run() -> Result<()> { &all_project_refs, &command_overrides, Some(&effective_primary_addrs), + true, ) } else { - executor.execute(&target, &all_project_refs, &graph, Some(&primary_addrs)) + executor.execute( + &target, + &all_project_refs, + &graph, + Some(&primary_addrs), + true, + ) }; // Output results based on mode @@ -1378,6 +1385,7 @@ fn run() -> Result<()> { &executor_projects, &command_overrides, Some(&primary_projects), + !run_args.no_deps, ) } else { // Pass ALL projects so executor can resolve target-level dependencies @@ -1394,6 +1402,7 @@ fn run() -> Result<()> { &executor_projects, &graph, Some(&primary_projects), + !run_args.no_deps, ) }; diff --git a/tests/dev_services.rs b/tests/dev_services.rs index fd9de4b..172d84c 100644 --- a/tests/dev_services.rs +++ b/tests/dev_services.rs @@ -302,6 +302,10 @@ fn exported_variable_targets_are_rejected_outside_services_up() { command = "sleep 30" stream = true exports_vars = true + +[targets.web] +command = "echo ok" +depends_on = ["//self:dev"] "#, ) .unwrap(); @@ -332,6 +336,7 @@ target = "//producer:dev" &["watch", "//producer:dev", "--no-initial"], &["run", "//consumer:dev"], &["dev", "//consumer"], + &["web", "//producer"], ]; for args in cases { let output = Command::new(env!("CARGO_BIN_EXE_aster")) @@ -351,17 +356,23 @@ target = "//producer:dev" ); } - let output = Command::new(env!("CARGO_BIN_EXE_aster")) - .args(["run", "//consumer:dev", "--no-deps"]) - .current_dir(root) - .output() - .unwrap(); - assert!( - output.status.success(), - "--no-deps should still run a non-exporter primary\nstdout={}\nstderr={}", - String::from_utf8_lossy(&output.stdout), - String::from_utf8_lossy(&output.stderr) - ); + for args in [ + ["run", "//consumer:dev", "--no-deps"].as_slice(), + ["dev", "//consumer", "--no-deps"].as_slice(), + ["web", "//producer", "--no-deps"].as_slice(), + ] { + let output = Command::new(env!("CARGO_BIN_EXE_aster")) + .args(args) + .current_dir(root) + .output() + .unwrap(); + assert!( + output.status.success(), + "{args:?} should still run a non-exporter primary\nstdout={}\nstderr={}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + } } #[test] From 7a55c2d305b0d2444f7797789d0b49ab624992f4 Mon Sep 17 00:00:00 2001 From: Arne Roomann-Kurrik Date: Tue, 15 Sep 2026 13:51:46 -0700 Subject: [PATCH 5/6] fix(services): fence exporter lifecycle transitions --- src/dev/runner.rs | 168 ++++++++++++++++++++++++++---- tests/dev_services.rs | 231 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 379 insertions(+), 20 deletions(-) diff --git a/src/dev/runner.rs b/src/dev/runner.rs index fe889d4..dd122b4 100644 --- a/src/dev/runner.rs +++ b/src/dev/runner.rs @@ -307,6 +307,7 @@ pub fn run_dev( &result.service, &plan, &mut runtimes, + &mut producers, &mut pending_starts, &system_tx, &mut dashboard, @@ -335,12 +336,31 @@ pub fn run_dev( generation, values, } => { + let Some(service) = plan + .services + .iter() + .find(|service| service.name == producer) + else { + continue; + }; + let Some(runtime) = runtimes.get(&producer) else { + continue; + }; + if runtime.generation != generation || runtime.process.is_none() { + continue; + } + let ExportResolution::Ready(current_inputs) = + resolve_exported_environment(service, &producers) + else { + continue; + }; + if !applied_exports_are_current(runtime, ¤t_inputs) { + continue; + } let Some(state) = producers.get_mut(&producer) else { continue; }; - if state.generation != generation - || runtimes[&producer].generation != generation - { + if state.generation != generation { continue; } state.epoch = state.epoch.saturating_add(1); @@ -372,28 +392,59 @@ pub fn run_dev( ExportResolution::Ready(_) if runtimes[&service.name].process.is_some() => { - queue_restart( - &mut pending_starts, - &service.name, - "effective exported environment changed", - &system_tx, - &mut dashboard, - ); - } - ExportResolution::Ready(_) | ExportResolution::Pending => {} - ExportResolution::Conflict(message) => { - emit_system(&system_tx, &service.name, message, true); - if runtimes[&service.name].process.is_some() { - hold_service( + if service.target.exports_vars() { + hold_exporter_tree( &service.name, - "exported variable sources conflict", + "effective exported environment changed", + &plan, &mut runtimes, + &mut producers, &mut pending_starts, &system_tx, &mut dashboard, #[cfg(unix)] session.as_ref(), ); + } else { + queue_restart( + &mut pending_starts, + &service.name, + "effective exported environment changed", + &system_tx, + &mut dashboard, + ); + } + } + ExportResolution::Ready(_) | ExportResolution::Pending => {} + ExportResolution::Conflict(message) => { + emit_system(&system_tx, &service.name, message, true); + if runtimes[&service.name].process.is_some() { + if service.target.exports_vars() { + hold_exporter_tree( + &service.name, + "exported variable sources conflict", + &plan, + &mut runtimes, + &mut producers, + &mut pending_starts, + &system_tx, + &mut dashboard, + #[cfg(unix)] + session.as_ref(), + ); + } else { + hold_service( + &service.name, + "exported variable sources conflict", + &mut runtimes, + &mut producers, + &mut pending_starts, + &system_tx, + &mut dashboard, + #[cfg(unix)] + session.as_ref(), + ); + } } } } @@ -425,6 +476,7 @@ pub fn run_dev( &producer, &plan, &mut runtimes, + &mut producers, &mut pending_starts, &system_tx, &mut dashboard, @@ -487,6 +539,7 @@ pub fn run_dev( &name, &plan, &mut runtimes, + &mut producers, &mut pending_starts, &system_tx, &mut dashboard, @@ -577,14 +630,13 @@ pub fn run_dev( true, ); if let Some(producer) = producers.get_mut(&service.name) { - let had_snapshot = producer.values.is_some(); producer.epoch = producer.epoch.saturating_add(1); producer.values = None; producer.healthy = false; producer.waiting_since = None; - if !had_snapshot && !supervisor_registered { + if !supervisor_registered { anyhow::bail!( - "variable exporter '{}' exited before publishing its initial snapshot", + "variable exporter '{}' exited before all variable exporters became ready", service.name ); } @@ -592,6 +644,7 @@ pub fn run_dev( &service.name, &plan, &mut runtimes, + &mut producers, &mut pending_starts, &system_tx, &mut dashboard, @@ -968,6 +1021,12 @@ fn resolve_exported_environment( }) } +fn applied_exports_are_current(runtime: &Runtime, resolved: &ResolvedExports) -> bool { + runtime.applied_exports == resolved.values + && runtime.applied_epochs == resolved.epochs + && runtime.applied_sources == resolved.sources +} + fn service_is_startable( service: &ServicePlan, runtimes: &HashMap, @@ -991,12 +1050,21 @@ fn hold_service( name: &str, reason: &str, runtimes: &mut HashMap, + producers: &mut HashMap, pending: &mut VecDeque<(String, String)>, system_tx: &std::sync::mpsc::Sender, dashboard: &mut Dashboard, #[cfg(unix)] session: Option<&SessionServer>, ) { let runtime = runtimes.get_mut(name).expect("runtime exists"); + if let Some(producer) = producers.get_mut(name) { + runtime.generation = runtime.generation.saturating_add(1); + producer.generation = runtime.generation; + producer.epoch = producer.epoch.saturating_add(1); + producer.values = None; + producer.healthy = false; + producer.waiting_since = None; + } #[cfg(unix)] runtime.export_endpoint.take(); if let Some(mut process) = runtime.process.take() { @@ -1015,11 +1083,48 @@ fn hold_service( queue_restart(pending, name, reason, system_tx, dashboard); } +#[allow(clippy::too_many_arguments)] +fn hold_exporter_tree( + name: &str, + reason: &str, + plan: &DevPlan, + runtimes: &mut HashMap, + producers: &mut HashMap, + pending: &mut VecDeque<(String, String)>, + system_tx: &std::sync::mpsc::Sender, + dashboard: &mut Dashboard, + #[cfg(unix)] session: Option<&SessionServer>, +) { + stop_dependent_services( + name, + plan, + runtimes, + producers, + pending, + system_tx, + dashboard, + #[cfg(unix)] + session, + ); + hold_service( + name, + reason, + runtimes, + producers, + pending, + system_tx, + dashboard, + #[cfg(unix)] + session, + ); +} + #[allow(clippy::too_many_arguments)] fn stop_dependent_services( producer: &str, plan: &DevPlan, runtimes: &mut HashMap, + producers: &mut HashMap, pending: &mut VecDeque<(String, String)>, system_tx: &std::sync::mpsc::Sender, dashboard: &mut Dashboard, @@ -1049,6 +1154,7 @@ fn stop_dependent_services( &service.name, "required exported variables unavailable", runtimes, + producers, pending, system_tx, dashboard, @@ -1684,6 +1790,28 @@ mod tests { )); } + #[test] + fn applied_export_epochs_fence_stale_nested_exporters() { + let runtime = Runtime { + process: None, + #[cfg(unix)] + export_endpoint: None, + generation: 1, + applied_exports: HashMap::from([("TOKEN".to_string(), "one".to_string())]), + applied_epochs: HashMap::from([("producer".to_string(), 3)]), + applied_sources: HashSet::from(["producer".to_string()]), + }; + let mut resolved = ResolvedExports { + values: runtime.applied_exports.clone(), + epochs: runtime.applied_epochs.clone(), + sources: runtime.applied_sources.clone(), + }; + assert!(applied_exports_are_current(&runtime, &resolved)); + + resolved.epochs.insert("producer".to_string(), 4); + assert!(!applied_exports_are_current(&runtime, &resolved)); + } + #[test] fn delayed_generated_event_is_ignored_but_changed_identity_restarts() { let temp = tempfile::tempdir().unwrap(); diff --git a/tests/dev_services.rs b/tests/dev_services.rs index 172d84c..d6699cf 100644 --- a/tests/dev_services.rs +++ b/tests/dev_services.rs @@ -288,6 +288,237 @@ stream = true terminate_aster(&mut aster); } +#[test] +fn nested_exporter_is_held_until_restarted_with_current_inputs() { + let temp = tempfile::tempdir().unwrap(); + let root = temp.path(); + fs::create_dir(root.join(".git")).unwrap(); + for project in ["producer", "nested", "blocker", "consumer"] { + fs::create_dir(root.join(project)).unwrap(); + fs::write( + root.join(project).join("package.json"), + format!(r#"{{"name":"{project}"}}"#), + ) + .unwrap(); + } + fs::write( + root.join("producer/publish.sh"), + r#"#!/bin/sh +printf '%s\n' '{"TOKEN":"one"}' > "$ASTER_EXPORT_VAR_PATH" +while [ ! -f ../nested.ready ] || [ ! -f ../blocker.started ]; do sleep 0.02; done +printf '%s\n' '{"TOKEN":"two"}' > "$ASTER_EXPORT_VAR_PATH" +sleep 30 +"#, + ) + .unwrap(); + fs::write( + root.join("nested/publish.sh"), + r#"#!/bin/sh +printf '{"DERIVED":"%s"}\n' "$TOKEN" > "$ASTER_EXPORT_VAR_PATH" +touch ../nested.ready +trap 'exit 0' TERM INT +while :; do sleep 1; done +"#, + ) + .unwrap(); + fs::write( + root.join("consumer/run.sh"), + r#"#!/bin/sh +printf 'consumer:%s\n' "$DERIVED" >> ../events.log +trap 'exit 0' TERM INT +while :; do sleep 1; done +"#, + ) + .unwrap(); + fs::write( + root.join("aster.toml"), + r#" +[dev.services.producer] +target = "//producer:dev" +order = 0 + +[dev.services.nested] +target = "//nested:dev" +inherit_env = ["TOKEN"] +order = 1 + +[dev.services.blocker] +target = "//blocker:dev" +order = 2 + +[dev.services.consumer] +target = "//consumer:dev" +inherit_env = ["DERIVED"] +order = 3 +"#, + ) + .unwrap(); + fs::write( + root.join("producer/aster.toml"), + r#" +[targets.dev] +command = "sh publish.sh" +stream = true +exports_vars = true +"#, + ) + .unwrap(); + fs::write( + root.join("nested/aster.toml"), + r#" +[targets.dev] +command = "sh publish.sh" +stream = true +exports_vars = true +depends_on = ["//producer:dev"] +"#, + ) + .unwrap(); + fs::write( + root.join("blocker/aster.toml"), + r#" +[targets.prepare] +command = "sh -c 'touch ../blocker.started; sleep 1'" + +[targets.dev] +command = "sleep 30" +stream = true +depends_on = ["//self:prepare"] +"#, + ) + .unwrap(); + fs::write( + root.join("consumer/aster.toml"), + r#" +[targets.dev] +command = "sh run.sh" +stream = true +depends_on = ["//nested:dev"] +"#, + ) + .unwrap(); + + let stdout = fs::File::create(root.join("stdout.log")).unwrap(); + let stderr = fs::File::create(root.join("stderr.log")).unwrap(); + let mut aster = Command::new(env!("CARGO_BIN_EXE_aster")) + .args(["services", "up", "--no-ui", "--no-watch"]) + .current_dir(root) + .stdout(stdout) + .stderr(stderr) + .spawn() + .unwrap(); + let events = root.join("events.log"); + if !condition_met(Duration::from_secs(15), || { + occurrences(&events, "consumer:two") == 1 + }) { + fail_with_process_diagnostics( + &mut aster, + &events, + &root.join("stdout.log"), + &root.join("stderr.log"), + "nested exporter did not restart with current inputs", + ); + } + assert_eq!(occurrences(&events, "consumer:one"), 0); + terminate_aster(&mut aster); +} + +#[test] +fn exporter_exit_before_global_readiness_fails_startup() { + let temp = tempfile::tempdir().unwrap(); + let root = temp.path(); + fs::create_dir(root.join(".git")).unwrap(); + for project in ["early", "slow"] { + fs::create_dir(root.join(project)).unwrap(); + fs::write( + root.join(project).join("package.json"), + format!(r#"{{"name":"{project}"}}"#), + ) + .unwrap(); + } + fs::write( + root.join("aster.toml"), + r#" +[dev.services.early] +target = "//early:dev" +order = 0 + +[dev.services.slow] +target = "//slow:dev" +order = 1 +"#, + ) + .unwrap(); + fs::write( + root.join("early/aster.toml"), + r#" +[targets.dev] +command = "sh publish.sh" +stream = true +exports_vars = true +"#, + ) + .unwrap(); + fs::write( + root.join("early/publish.sh"), + r#"#!/bin/sh +printf '%s\n' '{"EARLY":"ready"}' > "$ASTER_EXPORT_VAR_PATH" +sleep 0.3 +"#, + ) + .unwrap(); + fs::write( + root.join("slow/aster.toml"), + r#" +[targets.dev] +command = "sh publish.sh" +stream = true +exports_vars = true +"#, + ) + .unwrap(); + fs::write( + root.join("slow/publish.sh"), + r#"#!/bin/sh +sleep 3 +printf '%s\n' '{"SLOW":"ready"}' > "$ASTER_EXPORT_VAR_PATH" +sleep 30 +"#, + ) + .unwrap(); + + let stdout_path = root.join("stdout.log"); + let stderr_path = root.join("stderr.log"); + let stdout = fs::File::create(&stdout_path).unwrap(); + let stderr = fs::File::create(&stderr_path).unwrap(); + let mut aster = Command::new(env!("CARGO_BIN_EXE_aster")) + .args(["services", "up", "--no-ui", "--no-watch"]) + .current_dir(root) + .stdout(stdout) + .stderr(stderr) + .spawn() + .unwrap(); + + if !condition_met(Duration::from_secs(5), || { + aster.try_wait().unwrap().is_some() + }) { + fail_with_process_diagnostics( + &mut aster, + &root.join("events.log"), + &stdout_path, + &stderr_path, + "startup did not fail after a pre-ready exporter exited", + ); + } + let status = aster.try_wait().unwrap().expect("Aster exited"); + assert!(!status.success()); + let stderr = fs::read_to_string(stderr_path).unwrap_or_default(); + assert!( + stderr.contains("exited before all variable exporters became ready"), + "unexpected startup error: {stderr}" + ); +} + #[test] fn exported_variable_targets_are_rejected_outside_services_up() { let temp = tempfile::tempdir().unwrap(); From 42cf7efac382cb4507a638032911c3b2ab9475d1 Mon Sep 17 00:00:00 2001 From: Arne Roomann-Kurrik Date: Tue, 15 Sep 2026 16:27:11 -0700 Subject: [PATCH 6/6] fix(deps): update rustls for RUSTSEC-2026-0285 --- Cargo.lock | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 000de83..812bc94 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2112,9 +2112,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.43" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0283386ce02abc0151e1761d08802dfe86c173b0b494af5cbc086574e453da06" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "log", @@ -2145,9 +2145,9 @@ dependencies = [ [[package]] name = "rustls-webpki" -version = "0.103.13" +version = "0.103.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +checksum = "f3c3cf1d8b1e7d4927e2d154c3fcb02979afb9939629c62cd9048d4f07b60ac2" dependencies = [ "aws-lc-rs", "ring",