use crate::action::gate;
use crate::util::BErr;
use crate::work::tms;
use std::collections::BTreeMap;
use std::time::Duration;
const BUILD: &str = concat!(env!("CARGO_PKG_VERSION"), " (", env!("ROSTER_BUILD"), ")");
pub async fn run(cap: usize, once: bool, no_listen: bool, addr: Option<&str>) -> Result<(), BErr> {
if let Err(errors) = crate::config::load() {
for e in &errors {
eprintln!("config: {e}");
}
return Err(format!(
"invalid config ({} error(s)) — fix and retry, or: roster server validate",
errors.len()
)
.into());
}
if let Ok(c) = crate::config::snapshot() {
for w in &c.warnings {
eprintln!("warning: {w}");
}
}
let _daemon_lock =
match crate::statefile::FileLock::try_acquire_path(&crate::paths::lock_file("daemon")) {
Ok(Some(lock)) => lock,
Ok(None) => {
return Err(
"another roster server is already running for this deployment — \
stop it first, or check: roster server status"
.into(),
)
}
Err(e) => return Err(format!("could not take the daemon lock: {e}").into()),
};
bootstrap_llm_credential().await?;
if let Ok(c) = crate::config::snapshot() {
crate::run::boxed::pull_image(&c.box_image).await?;
}
let addrs = crate::gateway::resolve_bind_addrs(addr);
let gateways = crate::gateway::start(&addrs).await.map_err(|e| -> BErr {
if e.to_string().contains("Address already in use") {
format!(
"something is already listening on {} — another roster server? \
Check: roster server status (or pick a different --addr)",
addrs.join(", ")
)
.into()
} else {
e
}
})?;
eprintln!(
"roster server {BUILD} — gateway on {}; dispatch cap {cap}{}{}",
addrs.join(", "),
if once { "; once" } else { "" },
if no_listen { "; listeners off" } else { "" }
);
let mut listeners = Vec::new();
if !no_listen && !once {
let plan = crate::channel::listen::plan();
if plan.is_empty() {
eprintln!(
"listeners: none configured (a worker opts in via [channels] in its worker.toml)"
);
}
for (worker, platform, credential) in plan {
listeners.push(tokio::spawn(crate::channel::listen::supervised(
worker, platform, credential,
)));
}
if let Ok(c) = crate::config::snapshot() {
if let Some(first) = c.workers.first() {
eprintln!("talk to a worker from another terminal: roster talk {first}");
}
}
}
let dispatch = crate::work::dispatch::dispatch_loop(cap, once);
tokio::pin!(dispatch);
let result = tokio::select! {
r = &mut dispatch => r,
sig = shutdown_signal() => {
eprintln!("roster server: {sig} — shutting down");
tokio::time::sleep(Duration::from_secs(2)).await;
Ok(())
}
};
for g in gateways {
g.abort();
}
for l in listeners {
l.abort();
}
crate::gateway::clear_state();
result
}
async fn shutdown_signal() -> &'static str {
let mut term = match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
Ok(s) => s,
Err(_) => {
let _ = tokio::signal::ctrl_c().await;
return "SIGINT";
}
};
tokio::select! {
_ = term.recv() => "SIGTERM",
_ = tokio::signal::ctrl_c() => "SIGINT",
}
}
async fn bootstrap_llm_credential() -> Result<(), BErr> {
use crate::credential::LLM_PROVIDERS;
let present: Vec<&str> = LLM_PROVIDERS
.iter()
.copied()
.filter(|n| crate::credential::vault::get_credential(n).is_some())
.collect();
if !present.is_empty() {
for name in present {
crate::cli::connections::ensure_model_grant(name)?;
}
return Ok(());
}
let interactive = unsafe { libc::isatty(libc::STDIN_FILENO) } == 1;
if !interactive {
eprintln!(
"no LLM credential in the vault — boxes cannot call a model. \
Connect one: roster connection add anthropic (or openai-codex)"
);
return Ok(());
}
let pi_auth = std::path::PathBuf::from(std::env::var("HOME").unwrap_or_default())
.join(".pi/agent/auth.json");
let pi_logins = std::fs::read_to_string(&pi_auth)
.ok()
.and_then(|s| serde_json::from_str::<serde_json::Value>(&s).ok());
let mut imported = false;
if let Some(logins) = pi_logins.as_ref().and_then(|v| v.as_object()) {
for name in LLM_PROVIDERS {
let Some(cred) = logins.get(name).filter(|v| v.is_object()) else {
continue;
};
let answer = crate::credential::connect::ask(&format!(
"found a pi login for {name}; use it for roster? [y/N] "
))?;
if matches!(answer.trim(), "y" | "Y" | "yes") {
crate::credential::connect::store(name, cred)
.map_err(|e| format!("could not store {name}: {e}"))?;
eprintln!(
"imported {name} — roster now owns the token refresh; \
pi will re-login when it next needs to"
);
crate::cli::connections::ensure_model_grant(name)?;
imported = true;
}
}
}
if imported {
return Ok(());
}
let answer = crate::credential::connect::ask(
"no LLM credential yet — connect one now? [anthropic / openai-codex / skip] ",
)?;
match answer.trim() {
p @ ("anthropic" | "openai-codex") => {
crate::credential::connect::run(p)
.await
.map_err(|e| format!("connection add {p}: {e}"))?;
crate::cli::connections::ensure_model_grant(p)?;
}
_ => eprintln!("skipped — connect later with: roster connection add <provider>"),
}
Ok(())
}
pub fn validate() -> Result<(), BErr> {
match crate::config::load() {
Ok(c) => {
println!(
"config valid: {} worker(s) [{}], {} grant(s), {} action(s), {} trust rule(s), {} limit(s), {} heartbeat(s), {} listener(s), {} exposure(s)",
c.workers.len(),
c.workers.join(", "),
c.policy.rules.len(),
c.actions.actions.len(),
c.actions.trust.len(),
c.budget.limits.len(),
c.heartbeats.len(),
c.listeners.len(),
c.exposes.len(),
);
match &c.engine_dir {
Some(dir) if !dir.join("box").is_dir() => {
println!(
"warning: [engine] dir {} has no box/ — sessions will fail",
dir.display()
)
}
Some(dir) => println!(
"engine: dev override {} (mounted over the baked engine)",
dir.display()
),
None => println!("engine: baked into the box image ({})", c.box_image),
}
if !c.connections.is_empty() {
println!("connections: {}", c.connections.len());
}
for w in &c.warnings {
println!("warning: {w}");
}
Ok(())
}
Err(errors) => {
for e in &errors {
eprintln!("config: {e}");
}
Err(format!("{} error(s)", errors.len()).into())
}
}
}
pub enum GatewayHealth {
Down,
Up,
Foreign(String),
}
pub async fn gateway_health() -> (u16, GatewayHealth) {
let port = crate::gateway::recorded_port();
let resp = reqwest::Client::new()
.get(format!("http://127.0.0.1:{port}/healthz"))
.timeout(Duration::from_millis(700))
.send()
.await;
let health = match resp {
Ok(r) if r.status().is_success() => {
let v: serde_json::Value = r.json().await.unwrap_or_default();
match v.get("config_root").and_then(|s| s.as_str()) {
Some(root) if root == crate::paths::config_root().display().to_string() => {
GatewayHealth::Up
}
Some(root) => GatewayHealth::Foreign(root.to_string()),
None => GatewayHealth::Up,
}
}
_ => GatewayHealth::Down,
};
(port, health)
}
pub async fn gateway_up() -> bool {
matches!(gateway_health().await.1, GatewayHealth::Up)
}
pub async fn status(json: bool) -> Result<(), BErr> {
let (port, health) = gateway_health().await;
let gateway_up = matches!(health, GatewayHealth::Up);
let config = match crate::config::load() {
Ok(c) => format!("valid ({} worker(s))", c.workers.len()),
Err(errors) => format!(
"INVALID — {} error(s); run: roster server validate",
errors.len()
),
};
let mut queue_by_state: BTreeMap<String, usize> = BTreeMap::new();
for t in tms::list_all() {
*queue_by_state.entry(t.state).or_insert(0) += 1;
}
let gates_pending = gate::list_pending().len();
let listeners = crate::channel::listen::active_listeners();
if json {
let out = serde_json::json!({
"build": BUILD,
"gateway": {
"port": port,
"up": gateway_up,
"foreign_deployment": match &health {
GatewayHealth::Foreign(root) => Some(root.clone()),
_ => None,
},
},
"config": config,
"queue": queue_by_state,
"gates_pending": gates_pending,
"listeners": listeners.iter().map(|(worker, pid, since, alive)| serde_json::json!({
"worker": worker, "pid": pid, "since": since, "alive": alive,
})).collect::<Vec<_>>(),
});
println!("{}", serde_json::to_string_pretty(&out)?);
return Ok(());
}
println!("roster {BUILD}");
println!(
"gateway {}",
match &health {
GatewayHealth::Up => format!("up on :{port}"),
GatewayHealth::Down => format!("DOWN (nothing on :{port}) — run: roster server start"),
GatewayHealth::Foreign(root) => format!(
"DOWN for this deployment — :{port} is serving another deployment \
(config {root}); start this one on a free port: roster server start --addr 127.0.0.1:<port>"
),
}
);
println!("config {config}");
let queue_line = if queue_by_state.is_empty() {
"empty".to_string()
} else {
queue_by_state
.iter()
.map(|(state, n)| format!("{n} {state}"))
.collect::<Vec<_>>()
.join(", ")
};
println!(
"queue {queue_line}{}",
if !gateway_up && queue_by_state.get("pending").copied().unwrap_or(0) > 0 {
" (waiting for the server: roster server start)"
} else {
""
}
);
println!(
"gates {}",
if gates_pending == 0 {
"none pending".to_string()
} else {
format!("{gates_pending} PENDING — review: roster server approvals ls")
}
);
if listeners.is_empty() {
println!("listeners none");
} else {
for (worker, pid, since, alive) in listeners {
println!(
"listener {worker}: {} (pid {pid}, since {since})",
if alive {
"up"
} else {
"STALE LOCK — process gone"
}
);
}
}
Ok(())
}