use crate::auth::{
retrieve_credential, store_credential, CredentialStore, FileCredentialStore,
KeyringCredentialStore,
};
use crate::cli::commands::lifecycle;
use crate::cli::commands::proxy;
use crate::cli::output::{OutputConfig, OutputFormat};
use crate::cli::report::{Check, Report, Resolution, Section, SetupVerdict, State};
use crate::cli::ui;
use crate::cli::AuthLoginArgs;
use crate::cli::InitArgs;
use crate::config;
use crate::error::{OlError, ERR_INVALID_CONFIG, ERR_NO_CREDENTIALS, ERR_PORT_IN_USE};
use crate::hooks;
use crate::hooks::DetectedAgent;
use crate::telemetry::{self, config as telemetry_config, consent_file_path, Event};
use secrecy::ExposeSecret;
const STAGES: usize = 6;
const SETTLE_WAIT: std::time::Duration = std::time::Duration::from_secs(10);
pub fn run_init(args: &InitArgs, output: &OutputConfig) -> Result<(), OlError> {
crate::cli::header::print(output, &["init"]);
if args.dry_run {
return run_dry_run(args, output);
}
let _rail = ui::install(output, "init", STAGES);
ui::intro("Setting up this machine");
let result = run_stages(args, output);
if let Err(e) = &result {
report_not_live(e, output);
}
result
}
fn report_not_live(error: &OlError, output: &OutputConfig) {
if !ui::is_active() {
return;
}
if !ui::error_reported(error) {
output.print_error(error);
}
ui::outro("Stopped.");
let log = ui::log_path().map(|p| p.display().to_string());
let footer = log_footer(log.as_deref());
ui::card(&ui::Card {
verdict: ui::Verdict::NotLive,
title: "OpenLatch couldn't finish setting up".into(),
rows: vec![ui::Row::new("What happened", error.message.clone())],
sections: vec![
ui::Section::new(
"Try this",
vec![error
.suggestion
.clone()
.unwrap_or_else(|| "Run `openlatch init` again.".into())],
),
ui::Section::new(
"Need help",
vec!["Share the log below with your IT team or OpenLatch.".into()],
),
],
footer,
});
let mut fields = vec![("code", error.code)];
if let Some(log) = &log {
fields.push(("log", log));
}
ui::result(ui::Verdict::NotLive, error.exit_code(), &fields);
}
fn report_no_agent(output: &OutputConfig) {
if !ui::is_active() {
output.print_error(&hooks::agent_not_found_err());
return;
}
ui::halt("no AI agent found");
ui::outro("Paused: one thing to do first.");
ui::card(&ui::Card {
verdict: ui::Verdict::NoAgent,
title: "No AI agent to connect on this machine yet".into(),
rows: vec![
ui::Row::new(
"Looked for",
hooks::binding::DETECTABLE_AGENT_NAMES.join(" · "),
),
ui::Row::new("Status", "installed, not active yet"),
],
sections: vec![ui::Section::new(
"To finish",
vec!["Install one of them, then run `openlatch init` again.".into()],
)],
footer: footer_lines(),
});
ui::result(ui::Verdict::NoAgent, EXIT_NO_AGENT, &[]);
}
const EXIT_NO_AGENT: i32 = 3;
fn verdict_exit_code(verdict: SetupVerdict) -> i32 {
match verdict {
SetupVerdict::Live => 0,
SetupVerdict::NeedsAttention => crate::cli::report::EXIT_DEGRADED,
SetupVerdict::NotLive => 1,
}
}
struct CardFacts<'a> {
org: &'a str,
agents: &'a [(&'static str, &'static str)],
runs: &'a str,
no_start: bool,
console_url: &'a str,
}
fn final_checks(args: &InitArgs, output: &OutputConfig) -> Option<Report> {
let mut report = build_install_report(args, output)?;
let deadline = std::time::Instant::now() + SETTLE_WAIT;
while let Some(section) = settling(&report) {
if std::time::Instant::now() >= deadline {
break;
}
ui::progress(if section == Section::Policy {
"Waiting for first policies"
} else {
"Waiting for the model relay check"
});
std::thread::sleep(std::time::Duration::from_secs(1));
match build_install_report(args, output) {
Some(next) => report = next,
None => break,
}
}
Some(report)
}
fn settling(report: &Report) -> Option<Section> {
if report.setup_verdict() != SetupVerdict::Live {
return None;
}
report
.checks()
.iter()
.find(|c| c.state.is_warning() && c.resolution == Resolution::SelfResolving)
.map(|c| c.section)
}
fn report_verdict(report: &Report, facts: &CardFacts<'_>, exit_code: i32) {
let verdict = report.setup_verdict();
let failed: Vec<&Check> = report
.checks()
.iter()
.filter(|c| c.state.is_failure())
.collect();
match verdict {
SetupVerdict::Live => ui::done("all checks passed"),
SetupVerdict::NeedsAttention => ui::attention(&things(attention_items(report).len())),
SetupVerdict::NotLive => ui::failed(&format!(
"{} failed",
count(failed.len(), "check", "checks")
)),
}
ui::outro(if verdict == SetupVerdict::NotLive {
"Stopped."
} else {
"Done."
});
ui::card(&install_card(report, facts));
let agents = facts
.agents
.iter()
.map(|(wire, _)| *wire)
.collect::<Vec<_>>()
.join(",");
let log = ui::log_path().map(|p| p.display().to_string());
let mut fields: Vec<(&str, &str)> = [("org", facts.org), ("agents", agents.as_str())]
.into_iter()
.filter(|(_, v)| !v.is_empty())
.collect();
if verdict == SetupVerdict::NotLive {
if let Some(code) = failed.iter().find_map(|c| c.code) {
fields.push(("code", code));
}
if let Some(log) = &log {
fields.push(("log", log));
}
}
ui::result(card_verdict(verdict), exit_code, &fields);
}
fn card_verdict(verdict: SetupVerdict) -> ui::Verdict {
match verdict {
SetupVerdict::Live => ui::Verdict::Live,
SetupVerdict::NeedsAttention => ui::Verdict::NeedsAttention,
SetupVerdict::NotLive => ui::Verdict::NotLive,
}
}
fn attention_items(report: &Report) -> Vec<&Check> {
let needs_user = |c: &Check| c.state.is_warning() && c.resolution == Resolution::NeedsUser;
let blocking: Vec<Section> = report
.checks()
.iter()
.filter(|c| needs_user(c) && !matches!(c.state, State::Unknown(_)))
.map(|c| c.section)
.collect();
report
.checks()
.iter()
.filter(|c| match c.state {
State::Unknown(blocker) => needs_user(c) && !blocking.contains(&blocker),
_ => needs_user(c),
})
.collect()
}
fn count(n: usize, one: &str, many: &str) -> String {
format!("{n} {}", if n == 1 { one } else { many })
}
fn things(n: usize) -> String {
if n <= 1 {
"one thing left to do".to_string()
} else {
format!("{n} things left to do")
}
}
fn needs(n: usize) -> String {
if n <= 1 {
"one thing needs you".to_string()
} else {
format!("{n} things need you")
}
}
fn install_card(report: &Report, facts: &CardFacts<'_>) -> ui::Card {
let verdict = report.setup_verdict();
let mut rows = Vec::new();
if verdict != SetupVerdict::NotLive {
if !facts.org.is_empty() {
rows.push(ui::Row::new("Organization", facts.org));
}
let label = if facts.agents.len() == 1 {
"Agent"
} else {
"Agents"
};
for (i, (wire, name)) in facts.agents.iter().enumerate() {
rows.push(ui::Row::new(if i == 0 { label } else { "" }, *name));
rows.push(ui::Row::more(coverage_line(report, wire)));
}
}
let mut footer = footer_lines();
match verdict {
SetupVerdict::Live => {
rows.push(match report.section_state(Section::Policy) {
State::Ok => ui::Row::new("Policies", "in sync"),
State::Pending => ui::Row::new("Policies", "first sync in progress")
.with_tail("· applies on arrival"),
_ => ui::Row::new("Policies", "not in sync yet"),
});
rows.push(ui::Row::new("Runs", facts.runs));
let names: Vec<&str> = facts.agents.iter().map(|(_, n)| *n).collect();
ui::Card {
verdict: ui::Verdict::Live,
title: "OpenLatch is live on this machine".into(),
rows,
sections: vec![ui::Section::new(
"Next",
vec![
format!("Keep using {} as usual.", words(&names)),
format!("See what they do at `{}`", facts.console_url),
],
)],
footer,
}
}
SetupVerdict::NeedsAttention => {
let items = attention_items(report);
let mut next: Vec<String> = Vec::new();
for check in &items {
rows.push(ui::Row::new(check.section.title(), check.headline.clone()));
if let Some(remedy) = check.remedy.as_ref().filter(|r| !next.contains(r)) {
next.push(remedy.clone());
}
}
if next.is_empty() {
next.push("Run `openlatch doctor` for the details.".into());
}
let state = if facts.no_start { "set up" } else { "running" };
ui::Card {
verdict: ui::Verdict::NeedsAttention,
title: format!("OpenLatch is {state} — {}", needs(items.len())),
rows,
sections: vec![ui::Section::new("Next", next)],
footer,
}
}
SetupVerdict::NotLive => {
let mut fixes: Vec<String> = Vec::new();
for (i, check) in report
.checks()
.iter()
.filter(|c| c.state.is_failure())
.enumerate()
{
rows.push(ui::Row::new(
if i == 0 { "What happened" } else { "" },
format!("{}: {}", check.section.title(), check.headline),
));
if let Some(remedy) = check.remedy.as_ref().filter(|r| !fixes.contains(r)) {
fixes.push(remedy.clone());
}
}
if fixes.is_empty() {
fixes.push("Run `openlatch init` again.".into());
}
footer = log_footer(ui::log_path().map(|p| p.display().to_string()).as_deref());
ui::Card {
verdict: ui::Verdict::NotLive,
title: "OpenLatch couldn't finish setting up".into(),
rows,
sections: vec![
ui::Section::new("Try this", fixes),
ui::Section::new(
"Need help",
vec!["Share the log below with your IT team or OpenLatch.".into()],
),
],
footer,
}
}
}
}
fn coverage_line(report: &Report, agent: &str) -> String {
let actions = if report.agent_section_state(agent, Section::Hooks) == Some(State::Ok) {
"actions controlled"
} else {
"actions not controlled yet"
};
let calls = if report.agent_section_state(agent, Section::ModelRelay) == Some(State::Ok) {
"model calls observed"
} else {
"model calls not observed"
};
format!("{actions} · {calls}")
}
fn log_footer(log: Option<&str>) -> Vec<String> {
log.map(|log| vec![format!("Log {log}")])
.unwrap_or_default()
}
fn footer_lines() -> Vec<String> {
let mut footer =
vec!["Check anytime `openlatch status` · Full diagnostics `openlatch doctor`".into()];
if std::env::var_os("OPENLATCH_INSTALL_PATH_HINT").is_some_and(|v| !v.is_empty()) {
footer.push("Open a new terminal to use the `openlatch` command.".into());
}
footer
}
fn words(names: &[&str]) -> String {
match names {
[] => String::new(),
[one] => (*one).to_string(),
[head @ .., last] => format!("{} and {last}", head.join(", ")),
}
}
fn run_stages(args: &InitArgs, output: &OutputConfig) -> Result<(), OlError> {
ui::stage_named(
"Check this machine",
"Checking this machine",
"Checked this machine",
);
let mut ledger = InitLedger::default();
let ol_dir = config::openlatch_dir();
ledger.create_dir(&ol_dir).map_err(|e| {
OlError::new(
ERR_INVALID_CONFIG,
format!(
"Cannot create openlatch directory '{}': {e}",
ol_dir.display()
),
)
.with_suggestion("Check that you have write permission to your home directory.")
})?;
ledger.create_dir(&ol_dir.join("logs")).map_err(|e| {
OlError::new(
ERR_INVALID_CONFIG,
format!("Cannot create logs directory: {e}"),
)
})?;
let config_path = config::openlatch_dir().join("config.toml");
let re_init = config_path.exists();
let mut re_init_report = if re_init {
run_egress_gate(args, args.api_url.as_deref(), output)?
} else {
GateReport::default()
};
let reclaim = {
let pre_cfg = config::Config::load(None, None, false).unwrap_or_else(|_| {
let mut c = config::Config::defaults();
c.port = config::read_port_file().unwrap_or(c.port);
c
});
let outcome = lifecycle::reclaim_ports(&pre_cfg, output)?;
match outcome.action {
lifecycle::ReclaimAction::Nothing => output.print_step("No prior daemon to reclaim"),
lifecycle::ReclaimAction::Stopped => output.print_step(&format!(
"Reclaimed daemon ({})",
outcome.identity.describe()
)),
lifecycle::ReclaimAction::ForceKilled => output.print_step(&format!(
"Force-killed unresponsive daemon ({})",
outcome.identity.describe()
)),
}
outcome
};
let agents = hooks::detect_agents();
let agents = hooks::select_agents(agents, &args.agent)?;
if agents.is_empty() {
report_no_agent(output);
crate::cli::report::record_exit_code(EXIT_NO_AGENT);
return Ok(());
}
for a in &agents {
output.print_step(&format!("Detected agent: {}", agent_label(a)));
}
let found: Vec<&str> = agents.iter().map(DetectedAgent::display_name).collect();
let token_path = ol_dir.join("daemon.token");
let token_existed = token_path.exists();
let new_token = config::generate_token();
if !token_existed {
ledger.record_file(&token_path);
}
std::fs::write(&token_path, &new_token).map_err(|e| {
OlError::new(
ERR_INVALID_CONFIG,
format!("Cannot write token file '{}': {e}", token_path.display()),
)
.with_suggestion("Check that you have write permission to the openlatch directory.")
})?;
crate::fs_secure::restrict_to_owner(&token_path).map_err(|e| {
OlError::new(
ERR_INVALID_CONFIG,
format!("Cannot set permissions on token file: {e}"),
)
})?;
let token_action = if token_existed {
"(regenerated existing)"
} else {
"(new)"
};
output.print_step(&format!("Generated auth token {token_action}"));
let needs_port_probe = !config_path.exists() || args.reconfig;
let port = if needs_port_probe {
if args.reconfig {
let _ = lifecycle::run_stop(output);
let _ = std::fs::remove_file(config::openlatch_dir().join("daemon.port"));
let _ = std::fs::remove_file(&config_path);
}
let requested = std::env::var("OPENLATCH_PORT")
.ok()
.filter(|v| !v.trim().is_empty())
.map(|v| config::parse_port_env(&v))
.transpose()?;
let selected = match requested {
Some(p) => {
if std::net::TcpListener::bind(("127.0.0.1", p)).is_err() {
return Err(OlError::new(
ERR_PORT_IN_USE,
format!("OPENLATCH_PORT={p} is already in use"),
)
.with_suggestion(format!(
"Free port {p}, choose another via OPENLATCH_PORT, or unset it to probe {}-{} automatically.",
config::PORT_RANGE_START,
config::PORT_RANGE_END
))
.with_docs("https://docs.openlatch.ai/errors/OL-1500"));
}
output.print_substep(&format!("Selected port {p} (from OPENLATCH_PORT)"));
p
}
None => {
let probed =
config::probe_free_port(config::PORT_RANGE_START, config::PORT_RANGE_END)?;
output.print_substep(&format!("Selected port {probed} (first available)"));
probed
}
};
let config_existed = config_path.exists();
config::ensure_config(selected)?;
if !config_existed {
ledger.record_file(&config_path);
}
let port_file = config::openlatch_dir().join("daemon.port");
let port_file_existed = port_file.exists();
config::write_port_file(selected)?;
if !port_file_existed {
ledger.record_file(&port_file);
}
selected
} else {
config::Config::load(None, None, false)?.port
};
config::ensure_agent_id(&config_path)?;
if let Some(outcome) = re_init_report.deferred.take() {
proxy::persist_outcome(&config_path, &outcome, output)?;
}
if let Some(api_url) = &args.api_url {
config::persist_api_url(&config_path, api_url)?;
output.print_substep(&format!("Cloud API URL set to {api_url}"));
}
if args.no_model_relay {
config::persist_model_relay_enabled(&config_path, false)?;
output.print_substep(
"Model relay disabled in config — agents connect to the provider directly",
);
}
let cfg = config::Config::load(Some(port), None, false)?;
let gate_report = if re_init {
re_init_report
} else {
match run_egress_gate(args, args.api_url.as_deref(), output) {
Ok(report) => report,
Err(e) => {
ledger.unwind(output);
ui::take_back_log();
return Err(e);
}
}
};
ui::done(&format!("found {}", words(&found)));
ui::stage_named("Sign in", "Signing in", "Signed in");
let (auth_success, org_name) = run_auth_for_init(output, args.yes)?;
ui::done(&org_name);
ui::stage("Usage data", "Usage data");
let consent = handle_telemetry_consent(args, output, &ol_dir)?;
ui::done(consent);
let connecting: Vec<&str> = agents
.iter()
.filter(|a| a.installable())
.map(DetectedAgent::display_name)
.collect();
let (name, running) = match connecting.as_slice() {
[one] => (format!("Connect {one}"), format!("Connecting {one}")),
_ => (
"Connect your agents".to_string(),
"Connecting your agents".to_string(),
),
};
ui::stage_named(
&name,
&running,
&format!("Connected {}", words(&connecting)),
);
let staged = hooks::staging::stage_hook_binary(&ol_dir)?;
output.print_step(&format!("Hook binary at {}", staged.target().display()));
let installed = install_for_agents(&agents, cfg.port, &new_token, output)?;
let connected: Vec<&str> = installed.iter().map(|(a, _)| a.display_name()).collect();
if connected.len() < connecting.len() {
ui::relabel(&format!("Connected {}", words(&connected)));
}
ui::done("actions now go through OpenLatch");
ui::stage_named("Start OpenLatch", "Starting OpenLatch", "Started OpenLatch");
let (supervision_backend_label, supervision_mode_label, supervision_deferred_reason) =
run_supervision_install_for_init(args, &config_path, output);
let at_login = if supervision_mode_label == "active" {
"starts at login"
} else if args.no_persistence {
"won't restart at login (--no-persistence)"
} else {
"won't restart at login"
};
let start_plan = plan_daemon_start(
args.no_start,
args.foreground,
supervision_mode_label == "active",
);
let mut interrupted = false;
let mut runs = format!("in the background · {at_login}");
let (port, pid) = if start_plan == DaemonStartPlan::Skip {
output.print_step("Skipped daemon start (--no-start)");
runs = "not started (--no-start)".to_string();
(cfg.port, 0u32)
} else if start_plan == DaemonStartPlan::SupervisorOwned {
match lifecycle::verify_running_daemon(cfg.port, 10) {
Ok(pid) => {
output.print_step(&format!(
"Daemon started on port {} (PID {pid}, supervised, v{})",
cfg.port,
env!("OPENLATCH_VERSION")
));
(cfg.port, pid)
}
Err(lifecycle::StartFailure::VersionMismatch { serving, expected }) => {
let e = lifecycle::start_failure_error(
lifecycle::StartFailure::VersionMismatch { serving, expected },
cfg.port,
);
return Err(e);
}
Err(_)
if lifecycle::read_pid_file()
.filter(|p| lifecycle::is_process_alive(*p))
.is_some() =>
{
let pid = lifecycle::read_pid_file().unwrap_or(0);
output.print_step(&format!(
"Daemon starting under supervision (PID {pid}) — not yet answering /health"
));
output.print_info(
" Check `openlatch status` shortly, or the newest ~/.openlatch/logs/daemon.log.<date>.",
);
runs = format!("starting in the background · {at_login}");
(cfg.port, pid)
}
Err(_) => {
tracing::warn!(
"supervision reported active but no daemon came up; starting one directly"
);
start_and_prove(cfg.port, &new_token, output)?
}
}
} else if args.foreground {
output.print_notice(
"Run `openlatch doctor` in another shell to check this install while it is up",
);
output.print_step(&format!(
"Starting daemon on port {} (foreground)",
cfg.port
));
#[cfg(feature = "model-relay")]
let spawn_model_relay = cfg.model_relay.enabled;
#[cfg(not(feature = "model-relay"))]
let spawn_model_relay = false;
ui::done("running in this terminal");
ui::outro("Running in the foreground · Ctrl+C to stop");
interrupted = run_daemon_foreground(cfg.port, &new_token, spawn_model_relay)?
!= lifecycle::Stopped::Exited;
(cfg.port, std::process::id())
} else {
start_and_prove(cfg.port, &new_token, output)?
};
#[cfg(feature = "model-relay")]
verify_model_relay_preflight(args, &cfg, start_plan, output)?;
if auth_success {
let cloud_msg = if org_name.is_empty() {
"Cloud sync: enabled".to_string()
} else {
format!("Cloud sync: connected (org: {org_name})")
};
output.print_step(&cloud_msg);
output.print_info(" Events will be forwarded automatically");
}
if start_plan == DaemonStartPlan::Skip {
ui::relabel("Set up OpenLatch");
ui::done("not started (--no-start)");
} else {
ui::done(runs.replacen("in the background", "running", 1).as_str());
}
let on_rail = start_plan != DaemonStartPlan::Foreground;
if on_rail {
ui::stage_named("Final checks", "Running final checks", "Final checks");
}
let install_report = if interrupted {
None
} else {
final_checks(args, output)
};
let today = chrono::Local::now().format("%Y-%m-%d");
let log_path = config::openlatch_dir()
.join("logs")
.join(format!("events-{today}.jsonl"));
let verdict = install_report.as_ref().map(Report::setup_verdict);
let exit_code = verdict.map_or(0, verdict_exit_code);
if let Some(report) = &install_report {
crate::cli::report::record_exit_code(exit_code);
if on_rail {
let wired: Vec<(&'static str, &'static str)> = installed
.iter()
.map(|(a, _)| (a.agent_type(), a.display_name()))
.collect();
let api_url = cfg.cloud.api_url.clone();
let facts = CardFacts {
org: &org_name,
agents: &wired,
runs: &runs,
no_start: args.no_start,
console_url: &api_url,
};
report_verdict(report, &facts, exit_code);
}
}
for (a, result) in &installed {
telemetry::capture_global(Event::cli_initialized(
a.agent_type(),
result.entries.len(),
!token_existed,
));
}
telemetry::capture_global(Event::supervision_installed(
supervision_backend_label,
supervision_mode_label,
supervision_deferred_reason.as_deref(),
));
if output.format == OutputFormat::Json {
let agents_json: Vec<serde_json::Value> = installed
.iter()
.map(|(a, result)| {
serde_json::json!({
"agent": a.agent_type(),
"settings_path": a.settings_path().to_string_lossy(),
"entries": result.entries.len(),
"events": result
.entries
.iter()
.map(|e| e.event_type.as_str())
.collect::<Vec<_>>(),
})
})
.collect();
let cloud_status = if auth_success {
"connected"
} else {
"not_configured"
};
let json = serde_json::json!({
"status": install_report
.as_ref()
.map(|r| r.overall().key())
.unwrap_or("unknown"),
"exit_code": exit_code,
"verdict": verdict.map(SetupVerdict::key),
"reclaimed": {
"action": match reclaim.action {
lifecycle::ReclaimAction::Nothing => "nothing",
lifecycle::ReclaimAction::Stopped => "stopped",
lifecycle::ReclaimAction::ForceKilled => "force_killed",
},
"pid": reclaim.identity.pid,
"version": reclaim.identity.version,
"uptime_secs": reclaim.identity.uptime_secs,
"exe": reclaim.identity.exe,
},
"report": install_report.as_ref().map(crate::cli::report::Report::to_json),
"agents": agents_json,
"port": port,
"pid": pid,
"log_path": log_path.to_string_lossy(),
"token_action": token_action,
"daemon_started": !args.no_start,
"cloud_status": cloud_status,
"org_name": org_name,
"supervision": {
"mode": supervision_mode_label,
"backend": supervision_backend_label,
"disabled_reason": supervision_deferred_reason,
},
"proxy": {
"source": gate_report.source,
"url_masked": gate_report.url_masked,
"prompted": gate_report.prompted,
},
});
output.print_json(&json);
}
Ok(())
}
#[derive(Debug, Clone)]
enum InitArtifact {
File(std::path::PathBuf),
Directory(std::path::PathBuf),
}
#[derive(Default)]
struct InitLedger {
created: Vec<InitArtifact>,
}
impl InitLedger {
fn record_file(&mut self, path: &std::path::Path) {
self.created.push(InitArtifact::File(path.to_path_buf()));
}
fn create_dir(&mut self, path: &std::path::Path) -> std::io::Result<()> {
let existed = path.exists();
std::fs::create_dir_all(path)?;
if !existed {
self.created
.push(InitArtifact::Directory(path.to_path_buf()));
}
Ok(())
}
fn unwind(&self, output: &OutputConfig) {
if self.created.is_empty() {
return;
}
let mut dirs = Vec::new();
for artifact in &self.created {
match artifact {
InitArtifact::File(p) => {
let _ = std::fs::remove_file(p);
}
InitArtifact::Directory(p) => dirs.push(p),
}
}
dirs.sort_by_key(|p| std::cmp::Reverse(p.components().count()));
for dir in dirs {
let _ = std::fs::remove_dir(dir);
}
output.print_substep("Rolled back — this host is as `init` found it");
}
}
fn install_for_agents<'a>(
agents: &'a [DetectedAgent],
port: u16,
token: &str,
output: &OutputConfig,
) -> Result<Vec<(&'a DetectedAgent, hooks::HookInstallResult)>, OlError> {
let mut installed: Vec<(&DetectedAgent, hooks::HookInstallResult)> = Vec::new();
let mut install_failures: Vec<OlError> = Vec::new();
let mut skipped_non_installable: Vec<&'static str> = Vec::new();
for a in agents {
if !a.installable() {
output.print_notice(&format!(
"Skipping {} — detected, but this build writes no hooks into it.",
a.display_name()
));
skipped_non_installable.push(a.display_name());
continue;
}
match hooks::install_hooks(&*a.binding, port, token) {
Ok(result) => {
output.print_step(&format!("Hooks written to {}", a.settings_path().display()));
for entry in &result.entries {
let action_label = match entry.action {
hooks::HookAction::Added => "added",
hooks::HookAction::Replaced => "replaced",
};
output.print_substep(&format!("{} ({})", entry.event_type, action_label));
}
installed.push((a, result));
}
Err(e) => {
output.print_notice(&format!(
"Warning: could not write hooks for {}: {} ({})",
a.display_name(),
e.message,
e.code
));
install_failures.push(e);
}
}
}
if installed.is_empty() {
return Err(install_failures
.into_iter()
.next()
.unwrap_or_else(|| nothing_installed_err(&skipped_non_installable)));
}
Ok(installed)
}
fn nothing_installed_err(skipped_non_installable: &[&'static str]) -> OlError {
if skipped_non_installable.is_empty() {
return hooks::agent_not_found_err();
}
OlError::new(
crate::error::ERR_HOOK_AGENT_NOT_INSTALLABLE,
format!(
"No hooks were installed: {} detected, and this build writes no hooks into {}",
skipped_non_installable.join(", "),
if skipped_non_installable.len() == 1 {
"it"
} else {
"any of them"
}
),
)
.with_suggestion(
"Omit --agent to cover every agent on this host, or name one this build can install into."
.to_string(),
)
.with_docs("https://docs.openlatch.ai/errors/OL-1407")
}
fn run_dry_run(args: &InitArgs, output: &OutputConfig) -> Result<(), OlError> {
let config_path = config::openlatch_dir().join("config.toml");
let cfg = config::Config::load(None, None, false)?;
let api_url = args
.api_url
.clone()
.unwrap_or_else(|| cfg.cloud.api_url.clone());
let overrides = proxy::ProxyOverrides::from_init(args);
overrides.validate()?;
let mut egress_cfg = outgoing_egress(args, &cfg)?;
overrides.apply(&mut egress_cfg)?;
proxy::refuse_linux_pac(&egress_cfg)?;
let persisted_source = if args.reconfig {
None
} else {
proxy::PersistedProxy::read(&config_path).source
};
let probe_ok = probe_current_route(&api_url, &egress_cfg);
let (action, source) = if persisted_source.as_deref() == Some("manual") {
("keep", persisted_source.clone())
} else if probe_ok {
match &persisted_source {
Some(s) => ("keep", Some(s.clone())),
None if egress_cfg.has_proxy() => ("persist", Some("env".to_string())),
None => ("keep", None),
}
} else {
("discover", None)
};
if output.format == OutputFormat::Json {
let mut proxy_doc = serde_json::json!({ "action": action });
if let Some(s) = &source {
proxy_doc["source"] = serde_json::json!(s);
}
output.print_json(&serde_json::json!({
"status": "ok",
"dry_run": true,
"api_url": api_url,
"proxy": proxy_doc,
}));
} else {
output.print_step("Dry run — nothing was written");
output.print_substep(&format!("api_url {api_url}"));
match &source {
Some(s) => output.print_substep(&format!("proxy {action} (source: {s})")),
None => output.print_substep(&format!("proxy {action}")),
}
}
Ok(())
}
fn outgoing_egress(
args: &InitArgs,
cfg: &config::Config,
) -> Result<crate::egress::EgressConfig, OlError> {
if !args.reconfig {
return Ok(cfg.egress.clone());
}
crate::egress::EgressConfig::resolve(
None,
&crate::egress::ProcessEnv,
cfg.port,
cfg.model_relay.port,
)
}
fn probe_current_route(api_url: &str, egress_cfg: &crate::egress::EgressConfig) -> bool {
use crate::egress::CandidateProbe;
let Ok(probe) = crate::egress::HealthProbe::new(api_url, egress_cfg.clone()) else {
return false;
};
let via = egress_cfg
.url
.as_deref()
.filter(|_| egress_cfg.mode != crate::egress::ProxyMode::Direct)
.and_then(|u| reqwest::Url::parse(u).ok());
probe.probe(via.as_ref()).is_ok()
}
#[derive(Default)]
struct GateReport {
source: Option<String>,
url_masked: Option<String>,
prompted: bool,
deferred: Option<proxy::GateOutcome>,
}
fn run_egress_gate(
args: &InitArgs,
api_url_override: Option<&str>,
output: &OutputConfig,
) -> Result<GateReport, OlError> {
let cfg = config::Config::load(None, None, false)?;
let api_url = api_url_override
.map(str::to_string)
.unwrap_or_else(|| cfg.cloud.api_url.clone());
let overrides = proxy::ProxyOverrides::from_init(args);
overrides.validate()?;
let base = outgoing_egress(args, &cfg)?;
let mut terminal = crate::cli::prompt::TerminalPrompter::new(api_url.clone());
let prompter: Option<&mut dyn crate::cli::prompt::Prompter> =
if crate::cli::prompt::interactive(output, args.yes) {
Some(&mut terminal)
} else {
None
};
let config_path = config::openlatch_dir().join("config.toml");
let persisted = if args.reconfig {
proxy::PersistedProxy::default()
} else {
proxy::PersistedProxy::read(&config_path)
};
let outcome = match proxy::run_gate(&api_url, base, &overrides, &persisted, prompter, output) {
Ok(o) => o,
Err(failure) => {
report_gate_failure(&failure, &api_url, args, output);
return Err(failure.error);
}
};
let has_writes = !outcome.sets.is_empty() || !outcome.removes.is_empty();
if has_writes && config_path.exists() && !args.reconfig {
proxy::persist_outcome(&config_path, &outcome, output)?;
}
let probed = outcome.attempts.iter().filter(|a| a.was_probed()).count();
proxy::emit_proxy_configured(&outcome.config, probed, None);
let route = outcome
.config
.url
.as_deref()
.map(crate::egress::mask_userinfo);
match (outcome.source_str(), &route) {
(Some(source), Some(url)) => output.print_step(&format!("proxy via {url} ({source})")),
(Some(source), None) => {
output.print_step(&format!("Cloud reachable (proxy source: {source})"));
}
(None, _) => output.print_step("Cloud reachable (direct)"),
}
Ok(GateReport {
source: outcome.source_str().map(str::to_string),
url_masked: route,
prompted: outcome.prompted,
deferred: (has_writes && args.reconfig).then_some(outcome),
})
}
fn report_gate_failure(
failure: &proxy::GateFailure,
api_url: &str,
args: &InitArgs,
output: &OutputConfig,
) {
let headless = !crate::cli::prompt::interactive(output, args.yes);
let hint = if headless {
"No terminal to prompt on (or `--yes` was passed), so no proxy was requested. Set one \
with `openlatch init --proxy <url>` or `openlatch system proxy set <url>`, or run \
`openlatch init` interactively without `--yes`."
.to_string()
} else {
failure.error.suggestion.clone().unwrap_or_default()
};
if output.format == OutputFormat::Json {
output.print_json(&serde_json::json!({
"status": "failed",
"exit_code": 1,
"api_url": api_url,
"error": { "code": failure.error.code, "message": failure.error.message },
"message": hint,
"candidates": proxy::candidates_json(&failure.attempts),
}));
return;
}
output.print_error(&failure.error);
if headless {
output.print_note(&format!(" {hint}"));
}
for attempt in &failure.attempts {
if attempt.was_probed() {
output.print_note(&format!(" {}", attempt.trace_line()));
}
}
}
fn persist_env_key(
key: &str,
primary: &dyn CredentialStore,
fallback: &dyn CredentialStore,
output: &OutputConfig,
) {
let secret = secrecy::SecretString::from(key.to_string());
if let Err(e) = store_credential(primary, fallback, secret) {
output.print_notice(&format!(
"Could not store the API key for the daemon ({}): {}",
e.code, e.message
));
}
}
fn sign_in_needs_a_terminal(app_url: &str) -> OlError {
OlError::new(
ERR_NO_CREDENTIALS,
"no API key found, and no terminal to finish a browser sign-in on",
)
.with_suggestion(format!(
"Create an API key at {app_url} and re-run with it in OPENLATCH_API_KEY, or run \
`openlatch init` from a terminal without `--yes` to sign in through the browser."
))
}
fn run_auth_for_init(output: &OutputConfig, yes: bool) -> Result<(bool, String), OlError> {
let keyring = KeyringCredentialStore::new();
let cfg = config::Config::load(None, None, false).ok();
let agent_id = cfg
.as_ref()
.and_then(|c| c.agent_id.clone())
.unwrap_or_default();
let file_store =
FileCredentialStore::new(config::openlatch_dir().join("credentials.enc"), agent_id);
let api_url = cfg
.as_ref()
.map(|c| c.cloud.api_url.clone())
.unwrap_or_else(|| "https://app.openlatch.ai".to_string());
let egress = cfg
.as_ref()
.map(|c| c.egress.clone())
.unwrap_or_else(crate::egress::EgressConfig::direct);
let rt = tokio::runtime::Runtime::new().map_err(|e| {
OlError::new(
ERR_INVALID_CONFIG,
format!("Failed to create async runtime: {e}"),
)
})?;
if let Ok(val) = std::env::var("OPENLATCH_API_KEY") {
if !val.is_empty() {
persist_env_key(&val, &keyring, &file_store, output);
let (online, org_name, _org_id) = rt.block_on(
crate::cli::commands::auth::validate_online(&val, &api_url, &egress),
);
if online {
let msg = if org_name.is_empty() {
"Authenticated via env var".to_string()
} else {
format!("Authenticated via env var (org: {org_name})")
};
output.print_step(&msg);
return Ok((true, org_name));
}
output.print_step("Authenticated via env var (cloud offline - validation skipped)");
return Ok((true, String::new()));
}
}
if let Ok(existing_key) = retrieve_credential(&keyring, &file_store) {
let key_str = existing_key.expose_secret().to_string();
let v = rt.block_on(crate::cli::commands::auth::validate_online_full(
&key_str, &api_url, &egress,
));
if v.online {
let msg = if v.org_name.is_empty() {
"Authenticated".to_string()
} else {
format!("Authenticated (org: {})", v.org_name)
};
output.print_step(&msg);
return Ok((true, v.org_name));
}
if !v.rejected {
output.print_step("Authenticated (cloud offline - using stored credentials)");
return Ok((true, String::new()));
}
output.print_substep("Existing credentials invalid, re-authenticating...");
}
if !crate::cli::prompt::interactive(output, yes) {
return Err(sign_in_needs_a_terminal(
&crate::cli::commands::auth::resolve_app_url(),
));
}
let login_args = AuthLoginArgs { no_browser: false };
crate::cli::commands::auth::run_login(&login_args, output)?;
if let Ok(key) = retrieve_credential(&keyring, &file_store) {
let key_str = key.expose_secret().to_string();
let (_, org_name, _) = rt.block_on(crate::cli::commands::auth::validate_online(
&key_str, &api_url, &egress,
));
return Ok((true, org_name));
}
Ok((true, String::new()))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum DaemonStartPlan {
Skip,
Foreground,
SupervisorOwned,
SpawnBackground,
}
fn start_and_prove(port: u16, token: &str, output: &OutputConfig) -> Result<(u16, u32), OlError> {
let mut spawned = lifecycle::spawn_daemon_tracked(port, token)?;
match lifecycle::verify_started_daemon(&mut spawned, port, 10) {
Ok(()) => {
output.print_step(&format!(
"Daemon started on port {port} (PID {}, v{})",
spawned.pid,
env!("OPENLATCH_VERSION")
));
Ok((port, spawned.pid))
}
Err(failure) => Err(lifecycle::start_failure_error(failure, port)),
}
}
pub(crate) fn plan_daemon_start(
no_start: bool,
foreground: bool,
supervision_active: bool,
) -> DaemonStartPlan {
if no_start {
DaemonStartPlan::Skip
} else if foreground {
DaemonStartPlan::Foreground
} else if supervision_active {
DaemonStartPlan::SupervisorOwned
} else {
DaemonStartPlan::SpawnBackground
}
}
fn run_supervision_install_for_init(
args: &InitArgs,
config_path: &std::path::Path,
output: &OutputConfig,
) -> (&'static str, &'static str, Option<String>) {
use crate::supervision::{select_supervisor, SupervisionMode, SupervisorKind};
let skip_reason: Option<&'static str> = if args.foreground {
Some("foreground_session")
} else if args.no_start {
Some("no_start")
} else if args.no_persistence {
Some("user_opt_out")
} else if !crate::supervision::unreproducible_environment().is_empty() {
Some("isolated_instance")
} else {
None
};
if let Some(reason) = skip_reason {
let _ = config::persist_supervision_state(
config_path,
&SupervisionMode::Disabled,
&SupervisorKind::None,
Some(reason),
);
let msg = match reason {
"user_opt_out" => "Supervision: skipped (--no-persistence)",
"isolated_instance" => {
"Supervision: not applicable — an isolated instance is not machine-global"
}
"foreground_session" => "Supervision: skipped (foreground session)",
"no_start" => "Supervision: skipped (--no-start)",
_ => "Supervision: skipped",
};
output.print_step(msg);
return ("none", "disabled", Some(reason.to_string()));
}
let Some(supervisor) = select_supervisor() else {
let reason = "unsupported_os";
let _ = config::persist_supervision_state(
config_path,
&SupervisionMode::Deferred,
&SupervisorKind::None,
Some(reason),
);
output
.print_step("Supervision: deferred (no supported supervisor detected on this system)");
return ("none", "deferred", Some(reason.to_string()));
};
let exe_path =
std::env::current_exe().unwrap_or_else(|_| std::path::PathBuf::from("openlatch"));
let backend = supervisor.kind();
let backend_label: &'static str = match backend {
SupervisorKind::Launchd => "launchd",
SupervisorKind::Systemd => "systemd",
SupervisorKind::TaskScheduler => "task_scheduler",
SupervisorKind::None => "none",
};
match supervisor.install(&exe_path) {
Ok(()) => {
let _ = config::persist_supervision_state(
config_path,
&SupervisionMode::Active,
&backend,
None,
);
output.print_step(&format!(
"Supervision installed ({backend_label}) — daemon will auto-start on login"
));
output.print_info(
" Disable with `openlatch system supervision disable` or run `openlatch init --no-persistence`.",
);
(backend_label, "active", None)
}
Err(e) => {
let reason_text = format!("{} ({})", e.message, e.code);
let _ = config::persist_supervision_state(
config_path,
&SupervisionMode::Deferred,
&backend,
Some(&reason_text),
);
output.print_step(&format!(
"Supervision: deferred — {backend_label} install failed ({})",
e.code
));
output.print_info(&format!(" {}", e.message));
output.print_info(
" Init will continue; run `openlatch system supervision install` to retry after the issue is resolved.",
);
(backend_label, "deferred", Some(reason_text))
}
}
}
fn agent_label(agent: &DetectedAgent) -> String {
format!(
"{} ({})",
agent.display_name(),
agent.config_dir().display()
)
}
fn run_daemon_foreground(
port: u16,
token: &str,
spawn_model_relay: bool,
) -> Result<lifecycle::Stopped, OlError> {
let mut cfg = config::Config::load(Some(port), None, true)?;
cfg.foreground = true;
let rt = tokio::runtime::Runtime::new().map_err(|e| {
OlError::new(
ERR_INVALID_CONFIG,
format!("Failed to create async runtime: {e}"),
)
})?;
let token_owned = token.to_string();
let pid = std::process::id();
let supervised = lifecycle::supervisor_owns_this_daemon(&cfg.supervision);
#[cfg(feature = "crash-report")]
crate::telemetry::crash::set_daemon_scope(cfg.port, pid);
let restart_into = rt.block_on(async move {
use crate::envelope;
use crate::logging;
use crate::privacy;
let _guard = logging::daemon_log::init_daemon_logging(&cfg.log_dir, true);
if let Ok(deleted) = logging::cleanup_old_logs(&cfg.log_dir, cfg.retention_days) {
if deleted > 0 {
tracing::info!(deleted = deleted, "cleaned up old log files");
}
}
privacy::init_filter(&cfg.extra_patterns);
let pid_path = config::openlatch_dir().join("daemon.pid");
if let Err(e) = std::fs::write(&pid_path, pid.to_string()) {
tracing::warn!(error = %e, "failed to write PID file");
}
logging::daemon_log::log_startup(
env!("CARGO_PKG_VERSION"),
cfg.port,
pid,
envelope::os_string(),
envelope::arch_string(),
);
crate::cli::commands::lifecycle::log_observability_status_from_env();
let credential_store = crate::cli::commands::lifecycle::build_credential_store();
let restart_into = match crate::daemon::start_server(
cfg.clone(),
token_owned,
Some(credential_store),
spawn_model_relay,
)
.await
{
Ok(served) => {
eprintln!(
"openlatch daemon stopped \u{2022} uptime {} \u{2022} {} events processed",
crate::daemon::format_uptime(served.uptime_secs),
served.events
);
served.restart_into
}
Err(e) => {
tracing::error!(error = %e, "daemon exited with error");
eprintln!("Error: daemon exited unexpectedly: {e}");
None
}
};
let _ = std::fs::remove_file(&pid_path);
restart_into
});
if let Some(exe) = restart_into {
rt.shutdown_timeout(std::time::Duration::from_secs(5));
return lifecycle::hand_over(&exe, port, token, supervised);
}
Ok(lifecycle::Stopped::observe())
}
fn handle_telemetry_consent(
args: &InitArgs,
output: &OutputConfig,
ol_dir: &std::path::Path,
) -> Result<&'static str, OlError> {
let consent_path = consent_file_path(ol_dir);
if args.no_telemetry {
telemetry_config::write_consent(&consent_path, false)?;
output.print_step("Telemetry: disabled (--no-telemetry)");
return Ok("off (--no-telemetry)");
}
if args.telemetry {
telemetry_config::write_consent(&consent_path, true)?;
output.print_step("Telemetry: enabled (--telemetry)");
return Ok("sharing anonymous usage data (--telemetry)");
}
if consent_path.exists() {
let sharing = telemetry_config::read_consent(&consent_path)
.ok()
.flatten()
.is_some_and(|c| c.enabled);
return Ok(if sharing {
"unchanged (sharing)"
} else {
"unchanged (not sharing)"
});
}
let interactive = crate::cli::prompt::interactive(output, args.yes);
if !interactive {
telemetry_config::write_consent(&consent_path, false)?;
output.print_info(
"ℹ Telemetry is off in non-interactive mode. Enable with `openlatch system telemetry enable`.",
);
return Ok("off (no terminal to ask)");
}
let question = consent_question();
let outcome = ui::select(&question).unwrap_or_else(|| ask_without_rail(&question));
let (sharing, result) = consent_result(&outcome);
telemetry_config::write_consent(&consent_path, sharing)?;
if !matches!(outcome, ui::SelectOutcome::Answered(_)) {
ui::note("Turn it on anytime: `openlatch system telemetry enable`");
}
output.print_step(&format!("Usage data: {result}"));
Ok(result)
}
fn consent_question() -> ui::Select<'static> {
ui::Select {
title: "Usage data",
question: "Share anonymous usage data to help improve OpenLatch?",
lines: &[
"Command names, agent types, error codes and counts. Never prompts,",
"code, secrets, or anything that identifies you. Off anytime.",
],
yes: "Yes, share",
no: "No thanks",
default_yes: true,
timeout: Some(CONSENT_PROMPT_BUDGET),
}
}
fn ask_without_rail(question: &ui::Select<'_>) -> ui::SelectOutcome {
let caps = ui::Caps::new(ui::Mode::Plain, false, false, 80);
ui::Rail::new(caps, ui::Out::Term(console::Term::stderr()), None, 0).select(question)
}
fn consent_result(outcome: &ui::SelectOutcome) -> (bool, &'static str) {
use ui::SelectOutcome as O;
match outcome {
O::Answered(true) => (true, "sharing anonymous usage data"),
O::Answered(false) => (false, "not sharing"),
O::Eof => (false, "no answer (input closed) — usage data stays off"),
O::ReadError(_) => (false, "couldn't read your answer — usage data stays off"),
O::TimedOut => (false, "no answer in time — usage data stays off"),
O::Unrecognized => (false, "answer not recognized — usage data stays off"),
}
}
const CONSENT_PROMPT_BUDGET: std::time::Duration = std::time::Duration::from_secs(120);
#[cfg(feature = "model-relay")]
const PREFLIGHT_VERDICT_WAIT: std::time::Duration = std::time::Duration::from_secs(15);
#[cfg(feature = "model-relay")]
fn verify_model_relay_preflight(
args: &InitArgs,
cfg: &config::Config,
start_plan: DaemonStartPlan,
output: &OutputConfig,
) -> Result<(), OlError> {
if start_plan == DaemonStartPlan::Skip
|| args.foreground
|| !cfg.model_relay.enabled
|| args.no_model_relay
|| !cfg.model_relay.owns_agent_wiring()
{
return Ok(());
}
let planes: Vec<(&'static str, crate::model_relay::wire_format::WireFormat)> =
crate::hooks::detect_agents()
.iter()
.filter_map(|a| {
a.binding
.model_relay_wiring()
.map(|w| (a.agent_type(), w.wire_format))
})
.collect();
if !planes.is_empty() {
wait_for_request_planes(cfg, &planes, output)?;
}
wait_for_provider_endpoints(cfg, output);
Ok(())
}
#[cfg(feature = "model-relay")]
fn wait_for_request_planes(
cfg: &config::Config,
planes: &[(&'static str, crate::model_relay::wire_format::WireFormat)],
output: &OutputConfig,
) -> Result<(), OlError> {
let port = cfg.model_relay.port;
let deadline = std::time::Instant::now() + PREFLIGHT_VERDICT_WAIT;
let url = format!("http://127.0.0.1:{port}/admin/model-relay/status");
loop {
let status = crate::egress::blocking_client_builder()
.timeout(std::time::Duration::from_secs(2))
.build()
.ok()
.and_then(|c| c.get(&url).send().ok())
.and_then(|r| r.json::<serde_json::Value>().ok());
if let Some(body) = status {
match preflight_wait_rule(&body, planes, cfg) {
WaitRule::Ok => {
output.print_step(&format!(
"Model relay verified — agents routed via http://127.0.0.1:{port}"
));
return Ok(());
}
WaitRule::Err {
agent,
upstream,
reason,
} => {
let err = OlError::new(
crate::error::ERR_MODEL_RELAY_PREFLIGHT_FAILED,
format!("Model relay check failed for {agent}: {reason}"),
)
.with_suggestion(format!(
"{agent}'s request plane is absent: its settings were left untouched, so \
its sessions connect straight to the provider and keep working — but \
nothing is captured. Check network reachability to {upstream} (proxy, \
VPN, TLS interception), then run `openlatch restart`. Run `openlatch \
doctor` for the full picture."
))
.with_docs("https://docs.openlatch.ai/errors/OL-RELAY-PREFLIGHT");
return Err(err);
}
WaitRule::Wait => {}
}
}
if std::time::Instant::now() >= deadline {
output.print_substep(
"Model relay: no verdict yet — the daemon is still checking it. \
Run `openlatch doctor` in a moment to confirm.",
);
return Ok(());
}
std::thread::sleep(std::time::Duration::from_millis(250));
}
}
#[cfg(feature = "model-relay")]
fn wait_for_provider_endpoints(cfg: &config::Config, output: &OutputConfig) {
use crate::cli::commands::model_relay::{endpoint_rows_for, verify_endpoint_ownership};
let agents = crate::hooks::detect_agents();
if !agents
.iter()
.any(|a| a.binding.provider_endpoints().is_some())
{
return;
}
let deadline = std::time::Instant::now() + PREFLIGHT_VERDICT_WAIT;
loop {
if endpoints_settled(&endpoint_rows_for(cfg, &agents, &verify_endpoint_ownership)) {
return;
}
if std::time::Instant::now() >= deadline {
output.print_substep(
"Provider endpoints: the daemon is still wiring them. \
Run `openlatch doctor` in a moment to confirm.",
);
return;
}
std::thread::sleep(std::time::Duration::from_millis(500));
}
}
#[cfg(feature = "model-relay")]
fn endpoints_settled(rows: &[crate::cli::commands::model_relay::EndpointRow]) -> bool {
use crate::cli::commands::model_relay::EndpointState;
rows.iter()
.all(|r| !matches!(r.state, EndpointState::Unwired | EndpointState::Down))
}
#[cfg(feature = "model-relay")]
#[derive(Debug, PartialEq, Eq)]
enum WaitRule {
Ok,
Err {
agent: &'static str,
upstream: String,
reason: String,
},
Wait,
}
#[cfg(feature = "model-relay")]
fn preflight_wait_rule(
body: &serde_json::Value,
planes: &[(&'static str, crate::model_relay::wire_format::WireFormat)],
cfg: &config::Config,
) -> WaitRule {
let verdicts = body.get("preflight");
let mut all_ok = true;
for (agent, fmt) in planes {
match verdicts.and_then(|v| v.get(agent)).and_then(|v| v.as_str()) {
Some("ok") => {}
Some("failed") => {
return WaitRule::Err {
agent,
upstream: cfg.model_relay.upstream_for(*fmt),
reason: body
.get("preflight_error")
.and_then(|v| v.get(agent))
.and_then(|v| v.as_str())
.unwrap_or("the model relay could not complete a request to the provider")
.to_string(),
};
}
_ => all_ok = false,
}
}
if all_ok {
WaitRule::Ok
} else {
WaitRule::Wait
}
}
fn build_install_report(
args: &InitArgs,
output: &OutputConfig,
) -> Option<crate::cli::report::Report> {
let mut report = crate::cli::commands::doctor::run_all_checks(output)
.ok()?
.report;
if args.no_start {
report.replace_section(
Section::Daemon,
Check::off(Section::Daemon, "Not started at your request (--no-start)")
.code(crate::error::ERR_DAEMON_START_FAILED)
.source("--no-start")
.remedy("Run `openlatch start` when you want it up."),
);
for section in [
Section::Hooks,
Section::ModelRelay,
Section::Cloud,
Section::Policy,
Section::Inventory,
Section::Integrity,
] {
report.replace_section(section, Check::unknown(section, Section::Daemon));
}
}
Some(report)
}
#[cfg(test)]
mod start_plan_tests {
use super::*;
#[test]
fn active_supervision_means_init_does_not_spawn() {
assert_eq!(
plan_daemon_start(false, false, true),
DaemonStartPlan::SupervisorOwned
);
}
#[test]
fn without_supervision_init_still_spawns() {
assert_eq!(
plan_daemon_start(false, false, false),
DaemonStartPlan::SpawnBackground
);
}
#[test]
fn explicit_flags_win_over_supervision() {
assert_eq!(plan_daemon_start(true, false, true), DaemonStartPlan::Skip);
assert_eq!(
plan_daemon_start(false, true, true),
DaemonStartPlan::Foreground
);
assert_eq!(plan_daemon_start(true, true, true), DaemonStartPlan::Skip);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_headless_sign_in_names_the_key_and_where_to_get_one() {
let e = sign_in_needs_a_terminal("https://app.example.test");
assert_eq!(e.code, ERR_NO_CREDENTIALS);
let suggestion = e.suggestion.unwrap_or_default();
assert!(suggestion.contains("OPENLATCH_API_KEY"), "{suggestion}");
assert!(
suggestion.contains("https://app.example.test"),
"{suggestion}"
);
assert!(suggestion.contains("--yes"), "{suggestion}");
}
#[test]
fn consent_line_takes_the_printed_default() {
let default_yes = consent_question().default_yes;
for accepted in ["", "\n", " \r\n", "y", "Y", "yes", "YES"] {
assert_eq!(
ui::select::parse_answer(accepted, default_yes),
Some(true),
"{accepted:?} should take the [Y/n] default"
);
}
for declined in ["n", "N", "no", "NO", " no ", "No\r\n"] {
assert_eq!(
ui::select::parse_answer(declined, default_yes),
Some(false),
"{declined:?} should decline"
);
}
}
#[test]
fn only_an_actual_yes_enables_telemetry() {
use ui::SelectOutcome as O;
let failure = ui::select::ReadFailure {
kind: std::io::ErrorKind::InvalidInput,
os_code: Some(6),
};
assert!(consent_result(&O::Answered(true)).0);
for outcome in [
O::Answered(false),
O::Eof,
O::ReadError(failure),
O::TimedOut,
O::Unrecognized,
] {
assert!(!consent_result(&outcome).0, "{outcome:?} must not enable");
}
}
#[test]
fn consent_outcomes_keep_their_own_words() {
use ui::SelectOutcome as O;
let failure = ui::select::ReadFailure {
kind: std::io::ErrorKind::InvalidInput,
os_code: Some(6),
};
let words = |o: O| consent_result(&o).1;
assert_eq!(words(O::Answered(true)), "sharing anonymous usage data");
assert_eq!(words(O::Answered(false)), "not sharing");
assert!(words(O::ReadError(failure)).starts_with("couldn't read your answer"));
assert!(words(O::Eof).contains("input closed"));
assert!(words(O::TimedOut).contains("in time"));
let all = [O::Eof, O::ReadError(failure), O::TimedOut, O::Unrecognized];
for o in all {
let w = words(o);
assert!(w.ends_with("usage data stays off"), "{w}");
assert!(!w.contains("nothing answered"), "{w}");
}
}
#[test]
fn eof_is_not_the_printed_default() {
use ui::select::{ask_line, LineRead};
let default_yes = consent_question().default_yes;
let once = |read: LineRead| {
let mut read = Some(read);
ask_line(
default_yes,
None,
move |_| read.take().expect("one read"),
|| {},
)
};
assert_eq!(
once(LineRead::Line("\n".into())),
ui::SelectOutcome::Answered(true)
);
let eof = once(LineRead::Eof);
assert_eq!(eof, ui::SelectOutcome::Eof);
assert!(!consent_result(&eof).0);
}
#[cfg(feature = "model-relay")]
#[test]
fn init_waits_for_provider_slots_and_pending_is_settled() {
use crate::cli::commands::model_relay::{EndpointRow, EndpointState};
use crate::hooks::cline_providers::UncoveredReason;
let row = |state| EndpointRow {
agent: "cline",
key: "cline:gs:shared:ollamaBaseUrl".into(),
label: "ollama (ollamaBaseUrl)".into(),
file: None,
port: Some(7601),
state,
};
assert!(endpoints_settled(&[]));
assert!(endpoints_settled(&[
row(EndpointState::NextStart),
row(EndpointState::Active),
row(EndpointState::Uncovered(UncoveredReason::SignedHost)),
row(EndpointState::Verdict {
code: crate::error::ERR_MODEL_RELAY_PREFLIGHT_FAILED,
detail: "offline".into(),
}),
]));
assert!(!endpoints_settled(&[
row(EndpointState::NextStart),
row(EndpointState::Unwired)
]));
assert!(
!endpoints_settled(&[row(EndpointState::Down)]),
"a recorded port the daemon has not bound yet"
);
}
#[cfg(feature = "model-relay")]
#[test]
fn verify_model_relay_preflight_waits_for_every_agent() {
use crate::model_relay::wire_format::WireFormat;
let mut cfg = config::Config::defaults();
cfg.model_relay.upstream.insert(
WireFormat::OpenAiResponses.as_str().to_string(),
"http://127.0.0.1:9".to_string(),
);
let planes = [
("claude-code", WireFormat::AnthropicMessages),
("codex-cli", WireFormat::OpenAiResponses),
];
let pending = serde_json::json!({
"preflight": { "claude-code": "ok" },
"preflight_error": { "claude-code": null },
});
assert_eq!(
preflight_wait_rule(&pending, &planes, &cfg),
WaitRule::Wait,
"one green plane is not the host's answer while another has no verdict"
);
let still_pending = serde_json::json!({
"preflight": { "claude-code": "ok", "codex-cli": "pending" },
});
assert_eq!(
preflight_wait_rule(&still_pending, &planes, &cfg),
WaitRule::Wait
);
let failed = serde_json::json!({
"preflight": { "claude-code": "ok", "codex-cli": "failed" },
"preflight_error": { "claude-code": null, "codex-cli": "could not reach it" },
});
match preflight_wait_rule(&failed, &planes, &cfg) {
WaitRule::Err {
agent,
upstream,
reason,
} => {
assert_eq!(agent, "codex-cli");
assert_eq!(upstream, "http://127.0.0.1:9");
assert!(
!upstream.contains("anthropic"),
"the upstream named must be the failing format's, not Anthropic's"
);
assert_eq!(reason, "could not reach it");
}
other => panic!("a failed plane must fail the install, got {other:?}"),
}
let all_ok = serde_json::json!({
"preflight": { "claude-code": "ok", "codex-cli": "ok" },
});
assert_eq!(preflight_wait_rule(&all_ok, &planes, &cfg), WaitRule::Ok);
let legacy = serde_json::json!({ "preflight": "ok" });
assert_eq!(
preflight_wait_rule(&legacy, &planes, &cfg),
WaitRule::Wait,
"a per-agent reader must not accept a scalar verdict as everyone's"
);
}
use crate::hooks::binding::test_support::non_installable_agent;
fn silent_output() -> OutputConfig {
OutputConfig {
format: OutputFormat::Human,
verbose: false,
debug: false,
quiet: true,
color: false,
}
}
#[test]
fn init_skips_non_installable() {
let dir = tempfile::tempdir().expect("temp dir");
let agents = vec![non_installable_agent("cline", dir.path())];
let err = install_for_agents(&agents, 7590, "test-token", &silent_output())
.expect_err("nothing was installed, so the loop reports a failure");
assert_eq!(
err.code,
crate::error::ERR_HOOK_AGENT_NOT_INSTALLABLE,
"the skip is what failed the install, not a missing agent"
);
assert_eq!(
std::fs::read_dir(dir.path())
.expect("the asset root is readable")
.count(),
0,
"install_hooks was never reached: it creates the settings file, its \
parent, the HMAC key and the token, and none of them may exist"
);
}
#[test]
fn init_agent_cline_alone_refuses_with_its_own_code() {
let dir = tempfile::tempdir().expect("temp dir");
let agents = vec![non_installable_agent("cline", dir.path())];
let err = install_for_agents(&agents, 7590, "test-token", &silent_output())
.expect_err("a lone non-installable selection installs nothing");
assert_eq!(err.code, crate::error::ERR_HOOK_AGENT_NOT_INSTALLABLE);
assert_ne!(
err.code,
crate::error::ERR_HOOK_AGENT_NOT_FOUND,
"the agent WAS found — telling the operator to install it is wrong advice"
);
assert!(
err.message.contains("Cline"),
"the refusal must name the agent it refused: {}",
err.message
);
assert_eq!(
nothing_installed_err(&[]).code,
crate::error::ERR_HOOK_AGENT_NOT_FOUND,
"a host with no agent at all is still OL-1400 — the codes are two \
different facts, not one constant"
);
}
}
#[cfg(test)]
mod rail_tests {
use super::*;
use crate::cli::ui::{Caps, Mode, Rail, INSTALL_LOCK};
use std::sync::{Arc, Mutex};
fn lock() -> std::sync::MutexGuard<'static, ()> {
INSTALL_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn plain_rail() -> (Rail, Arc<Mutex<Vec<String>>>) {
Rail::buffer(Caps::new(Mode::Plain, false, false, 80), STAGES)
}
fn screen(s: &Arc<Mutex<Vec<String>>>) -> Vec<String> {
s.lock().unwrap().clone()
}
fn green() -> Report {
let mut report = Report::new();
for section in Section::ALL {
let check = Check::ok(section, "fine");
let check = if matches!(section, Section::Hooks | Section::ModelRelay) {
check.agent("claude-code")
} else {
check
};
report.push(check);
}
report
}
fn with(mut report: Report, check: Check) -> Report {
report.replace_section(check.section, check);
report
}
fn relay_uncovered() -> Check {
Check::off(Section::ModelRelay, "Cline's model calls are not relayed")
.code("OL-RELAY-UNCOVERED")
.remedy("None — a known coverage limit.")
.resolution(Resolution::NothingToDo)
}
fn first_bundle_pending() -> Check {
Check::pending(Section::Policy, "Waiting for the first policy bundle")
.code("OL-1210")
.remedy("Nothing to do — it applies on arrival.")
.resolution(Resolution::SelfResolving)
}
fn relay_off() -> Check {
Check::off(Section::ModelRelay, "Model relay is switched off in config")
.code("OL-RELAY-OFF")
.remedy("Run `openlatch system model-relay enable`.")
}
fn daemon_failed() -> Check {
Check::failed(
Section::Daemon,
"The service stopped right after it started",
)
.code("OL-1502")
.remedy("Run `openlatch init` again.")
}
const AGENTS: &[(&str, &str)] = &[("claude-code", "Claude Code")];
fn facts() -> CardFacts<'static> {
CardFacts {
org: "Acme Corp",
agents: AGENTS,
runs: "in the background · starts at login",
no_start: false,
console_url: "https://app.openlatch.ai",
}
}
#[test]
fn setup_verdict_maps_to_the_exit_table() {
assert_eq!(verdict_exit_code(SetupVerdict::Live), 0);
assert_eq!(
verdict_exit_code(SetupVerdict::NeedsAttention),
crate::cli::report::EXIT_DEGRADED
);
assert_eq!(verdict_exit_code(SetupVerdict::NotLive), 1);
assert_eq!(EXIT_NO_AGENT, 3, "no agent is `not found`, never a failure");
let limit = with(green(), relay_uncovered());
assert_eq!(limit.exit_code(), crate::cli::report::EXIT_DEGRADED);
assert_eq!(verdict_exit_code(limit.setup_verdict()), 0);
let off = with(green(), relay_off());
assert_eq!(verdict_exit_code(off.setup_verdict()), 7);
let broken = with(off, daemon_failed());
assert_eq!(verdict_exit_code(broken.setup_verdict()), 1);
}
#[test]
fn live_card_names_agents_their_coverage_and_policies() {
let report = with(green(), first_bundle_pending());
let card = install_card(&report, &facts());
assert_eq!(card.verdict, ui::Verdict::Live);
assert_eq!(card.title, "OpenLatch is live on this machine");
let rows: Vec<(&str, &str)> = card
.rows
.iter()
.map(|r| (r.label.as_str(), r.value.as_str()))
.collect();
assert_eq!(
rows,
vec![
("Organization", "Acme Corp"),
("Agent", "Claude Code"),
("", "actions controlled · model calls observed"),
("Policies", "first sync in progress"),
("Runs", "in the background · starts at login"),
]
);
assert_eq!(
card.sections[0].lines,
vec![
"Keep using Claude Code as usual.".to_string(),
"See what they do at `https://app.openlatch.ai`".to_string(),
]
);
assert!(card.footer[0].contains("`openlatch doctor`"));
}
#[test]
fn a_known_coverage_limit_is_a_coverage_line_not_a_task() {
let report = with(green(), relay_uncovered());
let card = install_card(&report, &facts());
assert_eq!(card.verdict, ui::Verdict::Live);
assert!(card
.rows
.iter()
.any(|r| r.value == "actions controlled · model calls not observed"));
}
#[test]
fn attention_card_lists_what_needs_the_user_and_its_remedy() {
let report = with(green(), relay_off());
let card = install_card(&report, &facts());
assert_eq!(card.verdict, ui::Verdict::NeedsAttention);
assert_eq!(card.title, "OpenLatch is running — one thing needs you");
assert!(card.rows.iter().any(
|r| r.label == "Model Relay" && r.value == "Model relay is switched off in config"
));
assert_eq!(
card.sections[0].lines,
vec!["Run `openlatch system model-relay enable`.".to_string()]
);
let mut no_start = facts();
no_start.no_start = true;
let card = install_card(&report, &no_start);
assert_eq!(card.title, "OpenLatch is set up — one thing needs you");
}
#[test]
fn not_live_card_says_what_happened_and_what_to_try() {
let report = with(green(), daemon_failed());
let card = install_card(&report, &facts());
assert_eq!(card.verdict, ui::Verdict::NotLive);
assert_eq!(card.title, "OpenLatch couldn't finish setting up");
assert_eq!(card.rows[0].label, "What happened");
assert_eq!(
card.rows[0].value,
"Daemon: The service stopped right after it started"
);
let labels: Vec<&str> = card.sections.iter().map(|s| s.label.as_str()).collect();
assert_eq!(labels, vec!["Try this", "Need help"]);
}
#[test]
fn an_unknown_waiting_on_a_healthy_section_is_still_named() {
let blocked = Check::unknown(Section::Policy, Section::Cloud)
.headline("policy_bundle compatibility unknown")
.code("OL-TEST")
.remedy("Fix Cloud first.");
let report = with(green(), blocked);
assert_eq!(report.setup_verdict(), SetupVerdict::NeedsAttention);
let items = attention_items(&report);
assert_eq!(items.len(), 1);
assert_eq!(items[0].section, Section::Policy);
let mut report = report;
report.replace_section(
Section::Cloud,
Check::off(Section::Cloud, "Cloud is off")
.code("OL-TEST")
.remedy("Turn it on."),
);
report.push(
Check::unknown(Section::Policy, Section::Cloud)
.code("OL-TEST")
.remedy("Fix Cloud first."),
);
let sections: Vec<Section> = attention_items(&report).iter().map(|c| c.section).collect();
assert_eq!(sections, vec![Section::Cloud]);
}
#[test]
fn settle_waits_only_on_a_self_resolving_check_of_an_otherwise_live_run() {
assert_eq!(settling(&green()), None);
assert_eq!(
settling(&with(green(), first_bundle_pending())),
Some(Section::Policy)
);
let amber = with(with(green(), first_bundle_pending()), relay_off());
assert_eq!(settling(&amber), None);
}
#[test]
fn plain_transcript_of_a_live_run() {
let _l = lock();
let (rail, s) = plain_rail();
{
let _guard = ui::install_rail(rail).expect("installs");
ui::intro("Setting up this machine");
ui::stage("Checking this machine", "Checked this machine");
ui::detail("Reclaimed daemon (PID 42)");
ui::done("found Claude Code");
ui::stage("Signing in", "Signed in");
ui::done("Acme Corp");
ui::stage("Usage data", "Usage data");
ui::done("not sharing");
ui::stage("Connecting Claude Code", "Connected Claude Code");
ui::done("actions now go through OpenLatch");
ui::stage("Starting OpenLatch", "Started OpenLatch");
ui::done("running · starts at login");
ui::stage("Running final checks", "Final checks");
report_verdict(&green(), &facts(), 0);
}
let lines = screen(&s);
assert_eq!(
lines[..9],
[
"Setting up this machine",
"[1/6] Checked this machine - found Claude Code",
"[2/6] Signed in - Acme Corp",
"[3/6] Usage data - not sharing",
"[4/6] Connected Claude Code - actions now go through OpenLatch",
"[5/6] Started OpenLatch - running - starts at login",
"[6/6] Final checks - all checks passed",
"Done.",
"",
]
);
assert_eq!(lines[9], " + OpenLatch is live on this machine");
assert_eq!(
lines.last().map(String::as_str),
Some("result=live org=\"Acme Corp\" agents=claude-code exit=0")
);
assert!(lines.iter().all(|l| l.is_ascii()), "{lines:#?}");
assert!(
!lines.iter().any(|l| l.contains("Reclaimed")),
"detail stays out of plain output without --verbose"
);
}
#[test]
fn plain_transcript_with_no_agent() {
let _l = lock();
let (rail, s) = plain_rail();
{
let _guard = ui::install_rail(rail).expect("installs");
ui::intro("Setting up this machine");
ui::stage("Checking this machine", "Checked this machine");
report_no_agent(&silent());
}
let lines = screen(&s);
assert_eq!(lines[1], "[1/6] Checked this machine - no AI agent found");
assert_eq!(lines[2], "Paused: one thing to do first.");
assert!(lines
.iter()
.any(|l| l.contains("No AI agent to connect on this machine yet")));
assert!(lines
.iter()
.any(|l| l.contains("Looked for Claude Code - Codex CLI - Cline")));
assert_eq!(
lines.last().map(String::as_str),
Some("result=attention reason=no-agent exit=3")
);
}
#[test]
fn plain_transcript_of_a_failed_stage() {
let _l = lock();
let (rail, s) = plain_rail();
let error = OlError::new("OL-1502", "The service stopped right after it started")
.with_suggestion("Run `openlatch init` again.");
{
let _guard = ui::install_rail(rail).expect("installs");
ui::stage("Starting OpenLatch", "Started OpenLatch");
report_not_live(&error, &silent());
}
let lines = screen(&s);
assert_eq!(
lines[..4],
[
"[1/6] Starting OpenLatch - failed",
" The service stopped right after it started (OL-1502)",
" Run openlatch init again.",
"Stopped.",
]
);
assert!(ui::error_reported(&error), "main must not print it again");
assert_eq!(
lines.last().map(String::as_str),
Some("result=not_live code=OL-1502 exit=1")
);
}
fn silent() -> OutputConfig {
OutputConfig {
format: OutputFormat::Human,
verbose: false,
debug: false,
quiet: false,
color: false,
}
}
#[test]
fn words_joins_names_the_way_a_sentence_does() {
assert_eq!(words(&[]), "");
assert_eq!(words(&["Cline"]), "Cline");
assert_eq!(words(&["Claude Code", "Cline"]), "Claude Code and Cline");
assert_eq!(
words(&["Claude Code", "Codex CLI", "Cline"]),
"Claude Code, Codex CLI and Cline"
);
}
}
#[cfg(test)]
mod env_key_tests {
use super::*;
use crate::auth::memory::InMemoryCredentialStore;
use secrecy::SecretString;
fn silent_output() -> OutputConfig {
OutputConfig {
format: OutputFormat::Human,
verbose: false,
debug: false,
quiet: true,
color: false,
}
}
struct RefusingStore;
impl CredentialStore for RefusingStore {
fn store(&self, _key: SecretString) -> Result<(), OlError> {
Err(OlError::new("OL-TEST", "store unavailable"))
}
fn retrieve(&self) -> Result<SecretString, OlError> {
Err(OlError::new("OL-TEST", "store unavailable"))
}
fn delete(&self) -> Result<(), OlError> {
Ok(())
}
}
#[test]
fn an_env_key_lands_in_the_store_browser_login_uses() {
let primary = InMemoryCredentialStore::new();
let fallback = InMemoryCredentialStore::new();
persist_env_key("olk_first", &primary, &fallback, &silent_output());
assert_eq!(primary.retrieve().unwrap().expose_secret(), "olk_first");
}
#[test]
fn a_re_init_with_a_different_key_replaces_the_stored_one() {
let primary = InMemoryCredentialStore::new();
let fallback = InMemoryCredentialStore::new();
persist_env_key("olk_first", &primary, &fallback, &silent_output());
persist_env_key("olk_second", &primary, &fallback, &silent_output());
assert_eq!(primary.retrieve().unwrap().expose_secret(), "olk_second");
}
#[test]
fn with_no_keychain_the_key_goes_to_the_fallback_store() {
let fallback = InMemoryCredentialStore::new();
persist_env_key("olk_first", &RefusingStore, &fallback, &silent_output());
assert_eq!(fallback.retrieve().unwrap().expose_secret(), "olk_first");
}
#[test]
fn a_failed_write_does_not_fail_init_or_echo_the_key() {
persist_env_key(
"olk_secret_value",
&RefusingStore,
&RefusingStore,
&silent_output(),
);
let e = store_credential(
&RefusingStore,
&RefusingStore,
SecretString::from("olk_secret_value".to_string()),
)
.unwrap_err();
assert!(!e.message.contains("olk_secret_value"), "{}", e.message);
}
}