use std::collections::HashMap;
use std::sync::Arc;
use areev_run::{CommandExecutor, ExecResult, HostToolExecutor};
use crate::flag;
pub fn flag_or_env(flags: &HashMap<String, String>, key: &str, var: &str) -> Option<String> {
flag(flags, key)
.or_else(|| std::env::var(var).ok())
.map(|v| v.trim().to_string())
.filter(|v| !v.is_empty())
}
pub fn recall_deadline(flags: &HashMap<String, String>) -> Result<Option<std::time::Duration>, String> {
match flag_or_env(flags, "recall-deadline-ms", "AREEV_RECALL_DEADLINE_MS") {
None => Ok(None),
Some(ms) => match ms.parse::<u64>() {
Ok(n) => Ok((n > 0).then(|| std::time::Duration::from_millis(n))),
Err(_) => Err(format!(
"--recall-deadline-ms (or $AREEV_RECALL_DEADLINE_MS) must be a whole number of \
milliseconds, got {ms:?}"
)),
},
}
}
pub fn host_facade(
m: areev_store::Areev,
ns: Option<String>,
flags: &HashMap<String, String>,
) -> Result<areev_cal::AreevFacade, String> {
let mut f = areev_cal::AreevFacade::with_session(m, ns, None);
f.set_recall_deadline(recall_deadline(flags)?);
Ok(f)
}
pub fn decider(
flags: &HashMap<String, String>,
m: &areev_store::Areev,
) -> Result<Option<std::sync::Arc<dyn areev_core::decide::DecisionBackend>>, String> {
match crate::resolve_decider(flags)? {
Some(chain) => Ok(Some(crate::decider_for_egress(chain, crate::store_egress_active(m))?)),
None => Ok(None),
}
}
pub struct NoExecutor;
impl HostToolExecutor for NoExecutor {
fn execute(
&self,
tool_name: &str,
_hash: &str,
_input: &serde_json::Value,
_idem: &str,
) -> ExecResult {
ExecResult::Err {
cause: areev_run_core::FailCause::ExecutorError,
detail: format!("no --tool-cmd configured; cannot execute host tool '{tool_name}'"),
}
}
}
pub fn egress_spec(flags: &HashMap<String, String>) -> Result<areev_run::EgressSpec, String> {
let ttl = match flag(flags, "credential-ttl") {
None => None,
Some(v) => Some(
v.trim()
.parse::<u64>()
.map_err(|_| format!("--credential-ttl: expected whole seconds, got {v:?}"))?,
),
};
areev_run::EgressSpec {
credentials: flag(flags, "credential").filter(|v| !v.trim().is_empty()),
allow_hosts: flag(flags, "allow-host").filter(|v| !v.trim().is_empty()),
tool_egress: flag(flags, "tool-egress").filter(|v| !v.trim().is_empty()),
credential_ttl_secs: ttl,
resolver_env: flag(flags, "resolver-env").filter(|v| !v.trim().is_empty()),
}
.with_env_fallback()
}
pub fn build_egress(
flags: &HashMap<String, String>,
) -> Result<Option<areev_run::Broker>, String> {
egress_spec(flags)?.build()
}
fn executor_timeout(flags: &HashMap<String, String>) -> Option<Option<std::time::Duration>> {
flag_or_env(flags, "executor-timeout", "AREEV_RUN_EXECUTOR_TIMEOUT")
.and_then(|v| v.parse::<u64>().ok())
.map(|secs| if secs == 0 { None } else { Some(std::time::Duration::from_secs(secs)) })
}
pub fn tool_env_policy(flags: &HashMap<String, String>) -> Option<areev_core::proc::EnvPolicy> {
let raw = match flag(flags, "tool-env") {
Some(v) => v,
None => std::env::var("AREEV_RUN_TOOL_ENV").ok()?,
};
let raw = raw.trim();
let (policy, dropped) = areev_run::env_allow_policy(if raw == "true" { "" } else { raw });
if !dropped.is_empty() {
eprintln!(
"areev: --tool-env dropped {} — already registered as holding a secret \
(--passphrase-env/--token-env/--credential). A tool never receives one.",
dropped.join(", ")
);
}
Some(policy)
}
pub fn connector_code(
flags: &HashMap<String, String>,
db: &str,
) -> Option<areev_trigger::ConnectorCode> {
let list = flag_or_env(flags, "allow-executor", "AREEV_RUN_ALLOW_EXECUTOR")?;
let allow: Vec<String> = list
.split(',')
.map(str::trim)
.filter(|a| !a.is_empty())
.map(str::to_string)
.collect();
if allow.is_empty() {
return None;
}
Some(areev_trigger::ConnectorCode {
allow,
cache_dir: flag_or_env(flags, "executor-cache", "AREEV_RUN_EXECUTOR_CACHE")
.map(std::path::PathBuf::from),
sandbox_cmd: flag_or_env(flags, "sandbox-cmd", "AREEV_RUN_SANDBOX_CMD"),
timeout: executor_timeout(flags),
env: tool_env_policy(flags),
db_locator: Some(db.to_string()),
})
}
pub fn tool_executor(
flags: &HashMap<String, String>,
egress: Option<&areev_run::EgressHandle>,
) -> Arc<dyn HostToolExecutor> {
let timeout = executor_timeout(flags);
let env = tool_env_policy(flags);
let base: Arc<dyn HostToolExecutor> = match flag_or_env(flags, "tool-cmd", "AREEV_RUN_TOOL_CMD")
{
Some(cmd) => {
let mut ce = CommandExecutor::new(&cmd);
if let Some(t) = timeout {
ce = ce.with_timeout(t);
}
if let Some(p) = env.clone() {
ce = ce.with_env_policy(p);
}
Arc::new(match egress {
Some(h) => ce.with_egress(h.clone()),
None => ce,
})
}
None => Arc::new(NoExecutor),
};
match flag_or_env(flags, "allow-executor", "AREEV_RUN_ALLOW_EXECUTOR") {
None => base,
Some(list) => {
let mut ce = areev_run::CodeExecutor::new(base);
for addr in list.split(',').map(str::trim).filter(|a| !a.is_empty()) {
ce = ce.allow(addr);
}
if let Some(dir) = flag_or_env(flags, "executor-cache", "AREEV_RUN_EXECUTOR_CACHE") {
ce = ce.cache_dir(dir);
}
if let Some(cmd) = flag_or_env(flags, "sandbox-cmd", "AREEV_RUN_SANDBOX_CMD") {
ce = ce.sandbox_cmd(&cmd);
}
if let Some(t) = timeout {
ce = ce.with_timeout(t);
}
if let Some(p) = env {
ce = ce.with_env_policy(p);
}
if let Some(h) = egress {
ce = ce.with_egress(h.clone());
}
Arc::new(ce)
}
}
}
pub fn can_execute(flags: &HashMap<String, String>) -> bool {
flag_or_env(flags, "tool-cmd", "AREEV_RUN_TOOL_CMD").is_some()
|| flag_or_env(flags, "allow-executor", "AREEV_RUN_ALLOW_EXECUTOR").is_some()
|| flag_or_env(flags, "model", "AREEV_RUN_MODEL").is_some()
}
pub fn toolcall_llm(
flags: &HashMap<String, String>,
) -> Result<Option<Arc<dyn areev_llm::ToolCallLlm>>, String> {
match flag_or_env(flags, "model", "AREEV_RUN_MODEL") {
None => Ok(None),
Some(spec) => areev_llm::resolve_toolcall(
&spec,
flag_or_env(flags, "base-url", "AREEV_RUN_BASE_URL").as_deref(),
flag_or_env(flags, "key-env", "AREEV_RUN_KEY_ENV").as_deref(),
)
.map(Some)
.map_err(|e| e.to_string()),
}
}
pub fn observer(
flags: &HashMap<String, String>,
) -> Result<Option<Arc<dyn areev_run::RunObserver>>, String> {
let mut observers: Vec<Arc<dyn areev_run::RunObserver>> = Vec::new();
if flag(flags, "events").is_some_and(|v| !matches!(v.as_str(), "false" | "0" | "off" | "no")) {
struct StderrEvents;
impl areev_run::RunObserver for StderrEvents {
fn event(&self, ev: &areev_run::RunEvent) {
if let Ok(line) = serde_json::to_string(ev) {
eprintln!("{line}");
}
}
}
observers.push(Arc::new(StderrEvents));
}
if let Some(endpoint) = flag(flags, "otel-endpoint") {
observers.push(Arc::new(areev_run::OtelObserver::new(&endpoint)?));
}
Ok(match observers.len() {
0 => None,
1 => observers.pop(),
_ => {
struct FanOut(Vec<Arc<dyn areev_run::RunObserver>>);
impl areev_run::RunObserver for FanOut {
fn event(&self, ev: &areev_run::RunEvent) {
for o in &self.0 {
o.event(ev);
}
}
}
Some(Arc::new(FanOut(observers)))
}
})
}
pub fn report_refusals(broker: &Option<Arc<areev_run::Broker>>) {
if let Some(b) = broker {
for r in b.refusals() {
eprintln!(
"areev: {} ({})",
areev_run_core::RunError::EgressRefused { destination: r.destination },
r.reason
);
}
}
}
pub fn run_options(flags: &HashMap<String, String>) -> areev_run::RunOptions {
areev_run::RunOptions {
budgets: areev_run::BudgetsSpec {
max_supersteps: flag(flags, "max-supersteps").and_then(|v| v.parse().ok()),
max_tokens: flag(flags, "max-tokens").and_then(|v| v.parse().ok()),
max_usd_micros: flag(flags, "max-usd")
.and_then(|v| v.parse::<f64>().ok())
.map(|usd| (usd * 1_000_000.0) as u64),
max_wall_ms: flag(flags, "max-wall-ms").and_then(|v| v.parse().ok()),
max_storage_bytes: flag(flags, "max-storage").and_then(|v| v.parse().ok()),
max_effects: flag(flags, "max-run-effects").and_then(|v| v.parse().ok()),
max_tool_calls: flag(flags, "max-tool-calls").and_then(|v| v.parse().ok()),
},
ask_ttl_sec: flag(flags, "ask-ttl").and_then(|v| v.parse().ok()),
workers: flag(flags, "workers").and_then(|v| v.parse().ok()).unwrap_or(4),
on_dangling: Default::default(),
llm_max_tokens: flag(flags, "llm-max-tokens").and_then(|v| v.parse().ok()),
max_effects_per_attempt: flag(flags, "max-effects").and_then(|v| v.parse().ok()),
llm_tool_result_chars: flag(flags, "llm-tool-result-chars")
.and_then(|v| v.parse().ok()),
llm_context_tokens: flag(flags, "llm-context-tokens").and_then(|v| v.parse().ok()),
inject_crash: None,
initiator: flag_or_env(flags, "initiator", "AREEV_RUN_INITIATOR"),
harness_ns: flag_or_env(flags, "harness-ns", "AREEV_RUN_HARNESS_NS"),
allow_confirmation_asks: flag_or_env(
flags,
"allow-confirmation-asks",
"AREEV_RUN_ALLOW_CONFIRMATION_ASKS",
)
.map(|v| v != "false" && v != "0")
.unwrap_or(false),
input_in_run_namespace: flag_or_env(
flags,
"input-placement",
"AREEV_RUN_INPUT_PLACEMENT",
)
.as_deref()
== Some("run-ns"),
lease_ms: flag_or_env(flags, "lease", "AREEV_RUN_LEASE")
.and_then(|v| v.parse::<i64>().ok())
.map(|secs| secs * 1000),
node: flag_or_env(flags, "node", "AREEV_NODE_ID"),
max_run_effects: flag_or_env(flags, "max-run-effects", "AREEV_RUN_MAX_RUN_EFFECTS")
.and_then(|v| v.parse().ok()),
max_tool_calls: flag_or_env(flags, "max-tool-calls", "AREEV_RUN_MAX_TOOL_CALLS")
.and_then(|v| v.parse().ok()),
max_concurrent_runs: flag_or_env(flags, "max-concurrent", "AREEV_RUN_MAX_CONCURRENT")
.and_then(|v| v.parse().ok()),
max_concurrent_runs_per_principal: flag_or_env(
flags,
"max-concurrent-per-principal",
"AREEV_RUN_MAX_CONCURRENT_PER_PRINCIPAL",
)
.and_then(|v| v.parse().ok()),
}
}
#[cfg(test)]
mod tests {
use super::*;
static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
fn env_guard() -> std::sync::MutexGuard<'static, ()> {
ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner())
}
fn flags(pairs: &[(&str, &str)]) -> HashMap<String, String> {
pairs.iter().map(|(k, v)| (k.to_string(), v.to_string())).collect()
}
#[test]
fn budget_flags_reach_the_run_a_firing_starts() {
let o = run_options(&flags(&[
("max-tokens", "5000"),
("max-usd", "0.25"),
("max-wall-ms", "60000"),
("ask-ttl", "3600"),
]));
assert_eq!(o.budgets.max_tokens, Some(5000));
assert_eq!(
o.budgets.max_usd_micros,
Some(250_000),
"--max-usd is dollars, stored as micros"
);
assert_eq!(o.budgets.max_wall_ms, Some(60_000));
assert_eq!(o.ask_ttl_sec, Some(3600));
}
#[test]
fn no_budget_flags_means_no_ceiling_not_a_surprise_one() {
let o = run_options(&flags(&[]));
assert_eq!(o.budgets.max_tokens, None);
assert_eq!(o.budgets.max_usd_micros, None);
assert_eq!(o.ask_ttl_sec, None);
assert_eq!(o.workers, 4, "the documented default");
assert_eq!(
o.max_effects_per_attempt, None,
"no flag = the manifest pins nothing = the run-core default"
);
assert_eq!(
o.llm_tool_result_chars, None,
"no flag = tool results enter the transcript verbatim, as they always did"
);
assert_eq!(
o.llm_context_tokens, None,
"no flag = no ceiling = no folds, as before the fold existed"
);
}
#[test]
fn the_effect_cap_reaches_the_run_a_firing_starts() {
let o = run_options(&flags(&[("max-effects", "40")]));
assert_eq!(o.max_effects_per_attempt, Some(40));
}
#[test]
fn the_tool_result_bound_reaches_the_run_a_firing_starts() {
let o = run_options(&flags(&[("llm-tool-result-chars", "12000")]));
assert_eq!(o.llm_tool_result_chars, Some(12_000));
}
#[test]
fn the_context_ceiling_reaches_the_run_a_firing_starts() {
let o = run_options(&flags(&[("llm-context-tokens", "150000")]));
assert_eq!(o.llm_context_tokens, Some(150_000));
}
#[test]
fn the_environment_stands_in_for_a_flag_a_heartbeat_cannot_carry() {
let var = "AREEV_TEST_STACK_SANDBOX";
std::env::set_var(var, "from-env");
assert_eq!(
flag_or_env(&flags(&[("sandbox-cmd", "from-flag")]), "sandbox-cmd", var).as_deref(),
Some("from-flag")
);
assert_eq!(flag_or_env(&flags(&[]), "sandbox-cmd", var).as_deref(), Some("from-env"));
std::env::set_var(var, " ");
assert_eq!(flag_or_env(&flags(&[]), "sandbox-cmd", var), None);
std::env::remove_var(var);
let addr = "1671652297b93a6a";
std::env::set_var("AREEV_RUN_ALLOW_EXECUTOR", addr);
let exec = tool_executor(&flags(&[]), None);
assert!(exec.code_allowed("tool-hash", &format!("cas://sha256:{addr}")));
std::env::remove_var("AREEV_RUN_ALLOW_EXECUTOR");
let exec = tool_executor(&flags(&[]), None);
assert!(!exec.code_allowed("tool-hash", &format!("cas://sha256:{addr}")));
}
#[test]
fn executor_timeout_parses_seconds_and_zero_means_wait_forever() {
assert_eq!(
executor_timeout(&flags(&[])),
None,
"unset means: the executor's own default stands"
);
assert_eq!(
executor_timeout(&flags(&[("executor-timeout", "45")])),
Some(Some(std::time::Duration::from_secs(45)))
);
assert_eq!(
executor_timeout(&flags(&[("executor-timeout", "0")])),
Some(None),
"0 restores the pre-1.3 wait-forever behaviour, on request rather than by omission"
);
assert_eq!(
executor_timeout(&flags(&[("executor-timeout", "not-a-number")])),
None,
"an unparseable value is ignored, like every other numeric flag run_options reads"
);
}
#[cfg(unix)]
#[test]
fn executor_timeout_flag_shortens_the_ceiling_a_tool_cmd_runs_under() {
let exec = tool_executor(
&flags(&[("tool-cmd", "while :; do :; done"), ("executor-timeout", "1")]),
None,
);
let started = std::time::Instant::now();
let result = exec.execute("work", "h", &serde_json::json!({}), "k");
assert!(
started.elapsed() < std::time::Duration::from_secs(10),
"the override must apply — the 300s default would still be sleeping"
);
match result {
ExecResult::Err { cause, detail } => {
assert_eq!(cause, areev_run::FailCause::Timeout, "{detail}");
}
ExecResult::Ok(v) => panic!("expected a timeout, got {v}"),
}
}
#[test]
fn tool_env_policy_reads_the_flag_and_treats_a_bare_one_as_clear_only() {
use areev_core::proc::EnvPolicy;
let _env = env_guard();
assert_eq!(tool_env_policy(&flags(&[])), None, "unset keeps the inherit default");
let minimal = EnvPolicy::minimal_allow();
match tool_env_policy(&flags(&[("tool-env", "true")])) {
Some(EnvPolicy::ClearExcept { allow }) => assert_eq!(allow, minimal),
other => panic!("a valueless --tool-env must still clear, got {other:?}"),
}
match tool_env_policy(&flags(&[("tool-env", "AWS_REGION, HTTPS_PROXY")])) {
Some(EnvPolicy::ClearExcept { allow }) => {
let mut want = minimal.clone();
want.extend(["AWS_REGION".to_string(), "HTTPS_PROXY".to_string()]);
assert_eq!(allow, want);
}
other => panic!("expected a cleared environment, got {other:?}"),
}
match tool_env_policy(&flags(&[("tool-env", "")])) {
Some(EnvPolicy::ClearExcept { allow }) => assert_eq!(allow, minimal),
other => panic!("an empty --tool-env must clear, not inherit, got {other:?}"),
}
match tool_env_policy(&flags(&[("tool-env", " ")])) {
Some(EnvPolicy::ClearExcept { allow }) => assert_eq!(allow, minimal),
other => panic!("a whitespace --tool-env must clear, not inherit, got {other:?}"),
}
}
#[test]
fn an_empty_tool_env_variable_still_clears() {
use areev_core::proc::EnvPolicy;
let _env = env_guard();
struct Restore(Option<String>);
impl Drop for Restore {
fn drop(&mut self) {
match self.0.take() {
Some(v) => std::env::set_var("AREEV_RUN_TOOL_ENV", v),
None => std::env::remove_var("AREEV_RUN_TOOL_ENV"),
}
}
}
let _restore = Restore(std::env::var("AREEV_RUN_TOOL_ENV").ok());
std::env::set_var("AREEV_RUN_TOOL_ENV", "");
match tool_env_policy(&flags(&[])) {
Some(EnvPolicy::ClearExcept { allow }) => {
assert_eq!(allow, EnvPolicy::minimal_allow());
}
other => panic!("AREEV_RUN_TOOL_ENV=\"\" must clear, not inherit, got {other:?}"),
}
std::env::remove_var("AREEV_RUN_TOOL_ENV");
assert_eq!(
tool_env_policy(&flags(&[])),
None,
"an ABSENT variable is what keeps the inherit default"
);
}
#[cfg(unix)]
#[test]
fn tool_env_clears_the_environment_and_passes_only_what_it_names() {
let _env = env_guard();
const PLANTED: &str = "AREEV_TEST_TOOL_ENV_PLANTED";
const NAMED: &str = "AREEV_TEST_TOOL_ENV_NAMED";
struct Planted;
impl Drop for Planted {
fn drop(&mut self) {
std::env::remove_var(PLANTED);
std::env::remove_var(NAMED);
}
}
let _planted = Planted;
std::env::set_var(PLANTED, "leaked");
std::env::set_var(NAMED, "kept");
let cmd = format!(
r#"printf '{{"planted":"%s","named":"%s","path":"%s"}}' "${PLANTED}" "${NAMED}" "${{PATH:+set}}""#
);
let seen = |extra: &[(&str, &str)]| {
let mut f = vec![("tool-cmd", cmd.as_str())];
f.extend_from_slice(extra);
match tool_executor(&flags(&f), None).execute("work", "h", &serde_json::json!({}), "k")
{
ExecResult::Ok(v) => v,
ExecResult::Err { detail, .. } => panic!("{detail}"),
}
};
let inherited = seen(&[]);
assert_eq!(inherited["planted"], "leaked", "the default still inherits");
let cleared = seen(&[("tool-env", NAMED)]);
assert_eq!(cleared["planted"], "", "an unnamed variable must not survive the clear");
assert_eq!(cleared["named"], "kept", "a named one must");
assert_eq!(cleared["path"], "set", "PATH is load-bearing — without it nothing resolves");
}
#[cfg(unix)]
#[test]
fn tool_env_refuses_to_re_admit_a_registered_secret() {
const SECRET: &str = "AREEV_TEST_TOOL_ENV_SECRET";
struct Planted;
impl Drop for Planted {
fn drop(&mut self) {
std::env::remove_var(SECRET);
}
}
let _planted = Planted;
std::env::set_var(SECRET, "hunter2");
areev_core::proc::deny_env_var(SECRET);
let cmd = format!(r#"printf '{{"seen":"%s"}}' "${SECRET}""#);
let out = tool_executor(
&flags(&[("tool-cmd", cmd.as_str()), ("tool-env", SECRET)]),
None,
)
.execute("work", "h", &serde_json::json!({}), "k");
match out {
ExecResult::Ok(v) => assert_eq!(
v["seen"], "",
"a registered secret named in --tool-env must still be withheld"
),
ExecResult::Err { detail, .. } => panic!("{detail}"),
}
}
}