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

Filter by extension

Filter by extension


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

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

48 changes: 47 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -376,7 +376,7 @@ If a project has a target named `services`, run it explicitly with
`aster target services <project selectors>`. 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.
Expand Down Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions src/cli/skills.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.<name>.proxy]`
with a distinct named `upstream_port` and proxy-specific `env`. `services up
<group> --proxy` keeps the service's advertised named port on the proxy and
Expand Down
52 changes: 52 additions & 0 deletions src/config/project.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<CacheConfig>,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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!(
Expand Down Expand Up @@ -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();
Expand Down
83 changes: 67 additions & 16 deletions src/dev/daemon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,10 @@
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;
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";
Expand Down Expand Up @@ -240,7 +241,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);
// 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);
const STOP_GRACE: Duration = Duration::from_secs(5);
Expand Down Expand Up @@ -271,6 +274,14 @@ mod platform {
services: Vec<String>,
ports: BTreeMap<String, u16>,
attach_socket: Option<PathBuf>,
#[serde(default = "ready_record_default")]
ready: bool,
#[serde(default)]
extend_deadline_secs: Option<u64>,
}

fn ready_record_default() -> bool {
true
}

pub(super) fn is_serve_invocation() -> bool {
Expand Down Expand Up @@ -419,14 +430,19 @@ mod platform {
Ok(stopped)
}

pub(super) fn register_ready(
workspace: &Path,
services: Vec<String>,
ports: BTreeMap<String, u16>,
) -> 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<String>, PathBuf)> {
let bundle_id = std::env::var(READY_ID_ENV).map_err(|_| {
daemon_error(DaemonErrorCode::InvalidRequest, "missing daemon bundle id")
})?;
Expand All @@ -436,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<String>,
ports: BTreeMap<String, u16>,
) -> 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::<PathBuf>::None,
"ready": false, "extend_deadline_secs": extra.as_secs(),
}))
}

pub(super) struct RuntimePaths {
Expand Down Expand Up @@ -1097,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;
Expand Down Expand Up @@ -1336,6 +1379,7 @@ mod platform {
#[cfg(not(unix))]
mod platform {
use super::*;
use std::time::Duration;
fn unsupported<T>() -> DaemonResult<T> {
Err(daemon_error(
DaemonErrorCode::UnsupportedPlatform,
Expand Down Expand Up @@ -1373,6 +1417,9 @@ mod platform {
) -> DaemonResult<()> {
Ok(())
}
pub(super) fn extend_ready_deadline(_: &Path, _: Duration) -> DaemonResult<()> {
Ok(())
}
}

pub fn ping_daemon() -> DaemonResult<u32> {
Expand Down Expand Up @@ -1414,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,
Expand Down
Loading
Loading