use crate::config::Config;
use crate::error::OlError;
use crate::hooks::bindings::cline::ClineBinding;
use crate::hooks::cline_app::{
self,
process::{self, ProcessTable},
ClineApp, HubState, RestartOutcome, RestartWhen, StaleWhy, HUB_STARTUP_POLL, IDLE_RECHECK,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum HubAction {
Nothing,
Start,
RestartIfIdle,
RestartNowCaptured,
RestartNowClean,
}
pub(super) fn hub_management_wanted(
owns_wiring: bool,
owns_machine: bool,
hub_relocated: bool,
) -> bool {
owns_wiring && (owns_machine || hub_relocated)
}
pub(crate) fn manages_hub(
config: &Config,
binding: &dyn crate::hooks::binding::AgentBinding,
) -> bool {
hub_management_wanted(
super::owns_wiring_for(config, binding),
crate::supervision::owns_machine_supervision(),
cline_app::hub_port() != cline_app::DEFAULT_HUB_PORT && !binding.config_is_machine_global(),
)
}
fn carries_capture(s: &HubState) -> bool {
matches!(
s,
HubState::Ours { .. }
| HubState::Stale {
why: StaleWhy::OlderApp { .. } | StaleWhy::ProxyChanged,
..
}
)
}
pub(super) fn decide(state: &HubState, live: bool) -> HubAction {
match (state, live) {
(HubState::Absent, true) => HubAction::Start,
(HubState::Absent, false) => HubAction::Nothing,
(HubState::Ours { .. }, true) => HubAction::Nothing,
(HubState::Foreign { .. }, _) => HubAction::Nothing,
(
HubState::Stale {
why: StaleWhy::NoProxyEnv,
..
},
true,
) => HubAction::RestartIfIdle,
(
HubState::Stale {
why: StaleWhy::NoProxyEnv,
..
},
false,
) => HubAction::Nothing,
(
HubState::Stale {
why: StaleWhy::OlderApp { .. },
..
},
true,
) => HubAction::RestartIfIdle,
(
HubState::Stale {
why: StaleWhy::ProxyChanged,
..
},
true,
) => HubAction::RestartNowCaptured,
(_, false) => HubAction::RestartNowClean,
}
}
pub(super) async fn hub_tick(
config: &Config,
relay_port: u16,
wiring: &crate::model_relay::preflight::WiringState,
agents: &[crate::hooks::DetectedAgent],
) {
let Some(cline) = agents
.iter()
.find(|a| a.agent_type() == ClineBinding::AGENT_TYPE)
else {
return;
};
if !manages_hub(config, &*cline.binding) {
return;
}
let url = format!("http://127.0.0.1:{relay_port}");
let live = wiring.is_wired(ClineBinding::AGENT_TYPE)
&& crate::hooks::proxy_env_file::live_proxy_url(ClineBinding::AGENT_TYPE).as_deref()
== Some(url.as_str());
let _ = tokio::task::spawn_blocking(move || tick_blocking(relay_port, live)).await;
}
fn tick_blocking(relay_port: u16, live: bool) {
let Some(app) = cline_app::detect() else {
return;
};
let port = cline_app::hub_port();
let state = cline_app::hub_state(port);
cline_app::forget_record_unless(&state);
let env = cline_app::ProxyEnvVars::for_relay(relay_port);
let r = match decide(&state, live) {
HubAction::Nothing => return,
HubAction::Start => {
cline_app::start_hub(&app, &env).map(|pid| RestartOutcome::Started { pid })
}
HubAction::RestartIfIdle => cline_app::restart_hub(&app, Some(&env), RestartWhen::IfIdle),
HubAction::RestartNowCaptured => cline_app::restart_hub(&app, Some(&env), RestartWhen::Now),
HubAction::RestartNowClean => cline_app::restart_hub(&app, None, RestartWhen::Now),
};
log(r);
}
fn log(r: Result<RestartOutcome, OlError>) {
match r {
Ok(RestartOutcome::Started { pid }) => tracing::debug!(pid, "hub start acknowledged"),
Ok(RestartOutcome::Restarted { old, new }) => {
tracing::info!(old, new = ?new, "Cline.app hub restarted");
}
Ok(RestartOutcome::DeferredBusy { clients }) => tracing::info!(
clients,
"Cline.app hub predates capture; quit Cline.app and keep it closed until the hub is \
recaptured (within a minute) — the hub outlives Cline.app"
),
Ok(RestartOutcome::IdleUnknown(w)) => tracing::info!(
reason = %w,
"Cline.app hub predates capture and its idleness is unknown; it is left running"
),
Ok(RestartOutcome::LeftForeign { pid }) => tracing::info!(
pid,
"Cline.app hub port held by a process OpenLatch does not manage"
),
Ok(RestartOutcome::NothingToDo) => {}
Err(e) => tracing::warn!(
code = %e.code,
error = %e.message,
"Cline.app hub could not be managed"
),
}
}
pub(crate) fn release_hub(config: &Config) {
let managed = crate::hooks::detect_agents()
.into_iter()
.find(|a| a.agent_type() == ClineBinding::AGENT_TYPE)
.map_or(
hub_management_wanted(
config.model_relay.owns_agent_wiring(),
crate::supervision::owns_machine_supervision(),
false,
),
|a| manages_hub(config, &*a.binding),
);
if !managed {
return;
}
release_with(
&*cline_app::process::table(),
cline_app::detect().as_ref(),
cline_app::hub_port(),
);
}
fn release_with(t: &dyn ProcessTable, app: Option<&ClineApp>, port: u16) {
let (state, info) = cline_app::hub_state_inspected(t, port, app);
if !carries_capture(&state) {
return;
}
match app {
Some(app) => log(cline_app::restart_hub_with(
t,
app,
None,
RestartWhen::Now,
port,
IDLE_RECHECK,
HUB_STARTUP_POLL,
)),
None => {
if let (HubState::Ours { pid } | HubState::Stale { pid, .. }, Some(info)) =
(state, info.as_ref())
{
let _ = process::terminate(t, pid, Some(&process::ProcIdentity::of(info)));
}
cline_app::remove_record();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::hooks::cline_app::test_fixture::{
clean_env, hub_fx, listen, our_env, record_exists, seed_record,
};
#[test]
fn hub_management_wanted_needs_wiring_and_a_hub_of_its_own() {
let rows = [
((true, true, true), true),
((true, true, false), true),
((true, false, true), true),
((true, false, false), false),
((false, true, true), false),
((false, true, false), false),
((false, false, true), false),
((false, false, false), false),
];
for ((w, m, r), want) in rows {
assert_eq!(
hub_management_wanted(w, m, r),
want,
"wiring={w} machine={m} relocated={r}"
);
}
}
#[test]
fn decide_covers_every_state() {
let stale = |why| HubState::Stale { pid: 5, why };
let older = || StaleWhy::OlderApp {
hub: "0.0.33".into(),
installed: "0.0.34".into(),
};
let rows = [
(HubState::Absent, true, HubAction::Start),
(HubState::Absent, false, HubAction::Nothing),
(HubState::Ours { pid: 5 }, true, HubAction::Nothing),
(HubState::Ours { pid: 5 }, false, HubAction::RestartNowClean),
(HubState::Foreign { pid: 5 }, true, HubAction::Nothing),
(HubState::Foreign { pid: 5 }, false, HubAction::Nothing),
(stale(StaleWhy::NoProxyEnv), true, HubAction::RestartIfIdle),
(stale(StaleWhy::NoProxyEnv), false, HubAction::Nothing),
(stale(older()), true, HubAction::RestartIfIdle),
(stale(older()), false, HubAction::RestartNowClean),
(
stale(StaleWhy::ProxyChanged),
true,
HubAction::RestartNowCaptured,
),
(
stale(StaleWhy::ProxyChanged),
false,
HubAction::RestartNowClean,
),
];
assert_eq!(rows.len(), 12);
for (state, live, want) in rows {
assert_eq!(decide(&state, live), want, "{state:?} live={live}");
}
}
#[test]
fn a_relocated_instance_never_reaches_the_process_table() {
let _state = crate::config::OPENLATCH_DIR_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let root = tempfile::tempdir().expect("tempdir");
let ol = root.path().join("openlatch");
std::fs::create_dir_all(&ol).expect("mkdir");
let apps = root.path().join("apps");
crate::hooks::cline_app::write_fake_bundle(
&apps,
crate::hooks::cline_app::BUNDLE_ID,
"0.0.34",
);
let _seam = crate::hooks::cline::cline_isolated([
("OPENLATCH_DIR", Some(ol.into_os_string())),
(
crate::hooks::cline::STORE_DIR_ENV,
Some(root.path().join("store").into_os_string()),
),
(
crate::hooks::cline::DATA_DIR_ENV,
Some(root.path().join("data").into_os_string()),
),
(
crate::hooks::cline::ASSETS_DIR_ENV,
Some(root.path().join("assets").into_os_string()),
),
(
crate::hooks::cline_app::APP_DIR_ENV,
Some(apps.into_os_string()),
),
(crate::hooks::cline_app::HUB_PORT_ENV, None),
]);
let _table = cline_app::process::install_for_tests(std::sync::Arc::new(
cline_app::process::test_support::PanickingTable,
));
assert!(
cline_app::detect().is_some(),
"the app is there to be managed"
);
assert!(!crate::supervision::owns_machine_supervision());
release_hub(&Config::load(None, None, false).expect("config"));
}
#[test]
fn release_takes_capture_off_an_owned_hub() {
let fx = hub_fx();
const RELAY: &str = "http://127.0.0.1:7600";
let t = fx.fake();
seed_record(&fx, 500, 1, "0.0.34", RELAY);
listen(&t, 500, fx.hub_info(Some(clean_env())));
*t.clients.lock().unwrap() = Ok(vec![500]);
release_with(&t, Some(&fx.app), fx.port);
assert_eq!(*t.terminated.lock().unwrap(), vec![500]);
assert!(!record_exists());
let t = fx.fake();
listen(&t, 500, fx.hub_info(Some(clean_env())));
release_with(&t, Some(&fx.app), fx.port);
assert!(t.terminated.lock().unwrap().is_empty());
assert!(t.spawns.lock().unwrap().is_empty());
let t = fx.fake();
seed_record(&fx, 500, 1, "0.0.34", RELAY);
listen(&t, 500, fx.hub_info(Some(clean_env())));
release_with(&t, None, fx.port);
assert_eq!(*t.terminated.lock().unwrap(), vec![500]);
assert!(t.spawns.lock().unwrap().is_empty());
assert!(!record_exists());
let t = fx.fake();
seed_record(&fx, 500, 1, "0.0.34", RELAY);
let mut info = fx.hub_info(Some(our_env(RELAY)));
info.start = Some(2);
listen(&t, 500, info);
assert_eq!(
cline_app::hub_state_with(&t, fx.port, None),
HubState::Foreign { pid: 500 }
);
release_with(&t, None, fx.port);
assert!(t.terminated.lock().unwrap().is_empty());
assert!(t.spawns.lock().unwrap().is_empty());
}
}