use crate::config::Config;
use crate::error::OlError;
use crate::hooks::bindings::cline::ClineBinding;
use crate::hooks::cline_app::{
self,
process::{self, ProcessTable},
ClineApp, Desired, HubState, RestartOutcome, RestartWhen, StaleWhy, HUB_STARTUP_POLL,
IDLE_RECHECK,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum HubAction {
Nothing,
Start,
RestartIfIdle,
RestartNow,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum HubOutcome {
AlreadyOk,
Started,
Restarted,
DeferredBusy,
IdleUnknown,
Failed(String),
NotManaged,
}
impl HubOutcome {
pub(crate) fn carries_desired(&self) -> bool {
matches!(self, Self::AlreadyOk | Self::Started | Self::Restarted)
}
fn of(r: &Result<RestartOutcome, OlError>) -> Self {
match r {
Ok(RestartOutcome::Started { .. }) => Self::Started,
Ok(RestartOutcome::Restarted { .. }) => Self::Restarted,
Ok(RestartOutcome::DeferredBusy { .. }) => Self::DeferredBusy,
Ok(RestartOutcome::IdleUnknown(_)) => Self::IdleUnknown,
Ok(RestartOutcome::LeftForeign { .. }) => Self::NotManaged,
Ok(RestartOutcome::NothingToDo) => Self::AlreadyOk,
Err(e) => Self::Failed(e.message.clone()),
}
}
}
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(),
)
}
pub(super) fn manages_hub_without_binding(config: &Config) -> bool {
hub_management_wanted(
config.model_relay.owns_agent_wiring(),
crate::supervision::owns_machine_supervision(),
false,
)
}
fn carries_ours(s: &HubState) -> bool {
matches!(
s,
HubState::Ours { .. }
| HubState::Stale {
why: StaleWhy::OlderApp { .. } | StaleWhy::ProxyChanged | StaleWhy::EnvChanged,
..
}
)
}
pub(crate) fn decide(state: &HubState, desired: &Desired) -> HubAction {
let empty = desired.is_empty();
match state {
HubState::Absent if empty => HubAction::Nothing,
HubState::Absent => HubAction::Start,
HubState::Ours { .. } => HubAction::Nothing,
HubState::Foreign { .. } => HubAction::Nothing,
HubState::Stale {
why: StaleWhy::NoOurEnv,
..
} if empty => HubAction::Nothing,
HubState::Stale {
why: StaleWhy::NoOurEnv,
..
} => HubAction::RestartIfIdle,
HubState::Stale {
why: StaleWhy::OlderApp { .. },
..
} => HubAction::RestartIfIdle,
HubState::Stale {
why: StaleWhy::ProxyChanged,
..
} => HubAction::RestartNow,
HubState::Stale {
why: StaleWhy::EnvChanged,
..
} => HubAction::RestartIfIdle,
}
}
pub(crate) fn reconcile(
t: &dyn ProcessTable,
app: &ClineApp,
port: u16,
desired: &Desired,
recheck: std::time::Duration,
poll: std::time::Duration,
) -> (HubAction, HubOutcome) {
let (state, info) = cline_app::hub_state_inspected(t, port, Some(app), desired);
cline_app::forget_record_unless(&state);
if let (HubState::Ours { pid }, Some(info)) = (&state, info.as_ref()) {
if let Err(e) = cline_app::adopt_record(*pid, info, app, port, desired) {
log(Err(e));
}
}
let action = decide(&state, desired);
let r = match action {
HubAction::Nothing => {
let outcome = if matches!(state, HubState::Foreign { .. }) {
HubOutcome::NotManaged
} else {
HubOutcome::AlreadyOk
};
return (action, outcome);
}
HubAction::Start => cline_app::start_hub_with(t, app, desired, port, poll)
.map(|pid| RestartOutcome::Started { pid }),
HubAction::RestartIfIdle => {
cline_app::restart_hub_with(t, app, desired, RestartWhen::IfIdle, port, recheck, poll)
}
HubAction::RestartNow => {
cline_app::restart_hub_with(t, app, desired, RestartWhen::Now, port, recheck, poll)
}
};
let outcome = HubOutcome::of(&r);
log(r);
(action, outcome)
}
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 does not carry OpenLatch's environment yet; quit Cline.app and keep it \
closed until the hub is restarted (within a minute) — the hub outlives Cline.app"
),
Ok(RestartOutcome::IdleUnknown(w)) => tracing::info!(
reason = %w,
"Cline.app hub does not carry OpenLatch's environment 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(manages_hub_without_binding(config), |a| {
manages_hub(config, &*a.binding)
});
if !managed {
return;
}
release_with(
&*cline_app::process::table(),
cline_app::detect().as_ref(),
cline_app::hub_port(),
IDLE_RECHECK,
HUB_STARTUP_POLL,
);
}
pub(super) fn release_with(
t: &dyn ProcessTable,
app: Option<&ClineApp>,
port: u16,
recheck: std::time::Duration,
poll: std::time::Duration,
) -> HubOutcome {
let empty = Desired::default();
let (state, info) = cline_app::hub_state_inspected(t, port, app, &empty);
if !carries_ours(&state) {
return HubOutcome::AlreadyOk;
}
match app {
Some(app) => {
let r =
cline_app::restart_hub_with(t, app, &empty, RestartWhen::Now, port, recheck, poll);
let outcome = HubOutcome::of(&r);
log(r);
outcome
}
None => {
let outcome = match (state, info.as_ref()) {
(HubState::Ours { pid } | HubState::Stale { pid, .. }, Some(info)) => {
match process::terminate(t, pid, Some(&process::ProcIdentity::of(info))) {
Ok(()) => HubOutcome::Restarted,
Err(e) => HubOutcome::Failed(e),
}
}
_ => HubOutcome::NotManaged,
};
cline_app::remove_record();
outcome
}
}
}
#[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_28() {
use crate::hooks::cline_app::{PluginEnv, ProxyEnvVars};
use HubAction::{Nothing, RestartIfIdle, RestartNow, Start};
let stale = |why| HubState::Stale { pid: 5, why };
let older = || StaleWhy::OlderApp {
hub: "0.0.33".into(),
installed: "0.0.34".into(),
};
let capture = ProxyEnvVars {
proxy_url: "http://127.0.0.1:7600".into(),
ca_pem: "/ol/ca.pem".into(),
no_proxy: "127.0.0.1:7600".into(),
};
let plugin = PluginEnv {
wrapper_path: "/ol/cline-plugin-bootstrap/wrapper".into(),
runtime: "/ol/bin/cline-js-runtime-hub".into(),
};
let desired = [
Desired::default(),
Desired {
capture: None,
plugin: Some(plugin.clone()),
},
Desired {
capture: Some(capture.clone()),
plugin: None,
},
Desired {
capture: Some(capture),
plugin: Some(plugin),
},
];
let table = [
(HubState::Absent, [Nothing, Start, Start, Start]),
(
HubState::Ours { pid: 5 },
[Nothing, Nothing, Nothing, Nothing],
),
(
HubState::Foreign { pid: 5 },
[Nothing, Nothing, Nothing, Nothing],
),
(
stale(StaleWhy::NoOurEnv),
[Nothing, RestartIfIdle, RestartIfIdle, RestartIfIdle],
),
(
stale(older()),
[RestartIfIdle, RestartIfIdle, RestartIfIdle, RestartIfIdle],
),
(
stale(StaleWhy::ProxyChanged),
[RestartNow, RestartNow, RestartNow, RestartNow],
),
(
stale(StaleWhy::EnvChanged),
[RestartIfIdle, RestartIfIdle, RestartIfIdle, RestartIfIdle],
),
];
let rows: Vec<(&HubState, &Desired, HubAction)> = table
.iter()
.flat_map(|(state, wants)| {
desired
.iter()
.zip(wants.iter())
.map(move |(d, want)| (state, d, *want))
})
.collect();
assert_eq!(rows.len(), 28);
for (state, d, want) in rows {
assert_eq!(decide(state, d), want, "{state:?} desired={d:?}");
}
}
#[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_plugin_off_an_owned_hub() {
use crate::hooks::cline_app::{test_fixture::seed_record_for, PluginEnv, ProxyEnvVars};
let fx = hub_fx();
const RELAY: &str = "http://127.0.0.1:7600";
let plugin = PluginEnv {
wrapper_path: fx
.root
.path()
.join("cline-plugin-bootstrap")
.join("wrapper"),
runtime: fx.root.path().join("bin").join("cline-js-runtime-hub"),
};
for recorded in [
Desired {
capture: Some(ProxyEnvVars::for_relay(7600)),
plugin: Some(plugin.clone()),
},
Desired {
capture: None,
plugin: Some(plugin.clone()),
},
] {
let t = fx.fake();
seed_record_for(&fx, 500, 1, "0.0.34", &recorded);
listen(&t, 500, fx.hub_info(Some(clean_env())));
*t.clients.lock().unwrap() = Ok(vec![500, 4242]);
release_with(&t, Some(&fx.app), fx.port, IDLE_RECHECK, HUB_STARTUP_POLL);
assert_eq!(*t.terminated.lock().unwrap(), vec![500], "{recorded:?}");
let spawns = t.spawns.lock().unwrap().clone();
assert_eq!(spawns.len(), 1, "{recorded:?}");
assert_eq!(
spawns[0].3,
cline_app::hub_env(fx.port, &Desired::default()),
"{recorded:?}: the release spawn carries the hub-mode env only"
);
assert!(
!spawns[0].3.iter().any(|(k, _)| {
crate::core::login_env::ALLOWED.contains(&k.as_str())
|| k.eq_ignore_ascii_case("HTTPS_PROXY")
}),
"{recorded:?}"
);
assert!(
!record_exists(),
"{recorded:?}: the clean spawn is unrecorded"
);
}
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, IDLE_RECHECK, HUB_STARTUP_POLL);
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, IDLE_RECHECK, HUB_STARTUP_POLL);
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, IDLE_RECHECK, HUB_STARTUP_POLL);
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, &Desired::default()),
HubState::Foreign { pid: 500 }
);
release_with(&t, None, fx.port, IDLE_RECHECK, HUB_STARTUP_POLL);
assert!(t.terminated.lock().unwrap().is_empty());
assert!(t.spawns.lock().unwrap().is_empty());
}
}