use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::{watch, Notify};
use crate::config::Config;
use crate::core::login_env::{self, Ledger, LoginEnvOutcome, Sink};
use crate::core::supervision::task::{spawn_supervised, HealthRegistry, RestartPolicy, TaskSpec};
use crate::hooks::bindings::cline::ClineBinding;
use crate::hooks::cline_app::{
self, process::ProcessTable, ClineApp, Desired, PluginEnv, ProxyEnvVars,
};
use crate::hooks::cline_hosts::HostClass;
use crate::hooks::cline_runtime::{self, reason, HostAssessment};
use super::cline_hub::{self, HubAction, HubOutcome};
pub(crate) const TICK: Duration = Duration::from_secs(60);
pub(crate) const STATE_FILE: &str = "cline-plugin-delivery.json";
pub(crate) type CaptureState = Option<Option<ProxyEnvVars>>;
#[derive(Clone)]
pub(crate) struct CapturePublisher {
tx: Arc<watch::Sender<CaptureState>>,
nudge: Arc<Notify>,
}
impl CapturePublisher {
pub(crate) fn publish(&self, capture: Option<ProxyEnvVars>) {
let next = Some(capture);
let changed = self.tx.send_if_modified(|current| {
if *current == next {
false
} else {
*current = next;
true
}
});
if changed {
self.nudge.notify_one();
}
}
}
pub(crate) struct CaptureWatch {
rx: watch::Receiver<CaptureState>,
nudge: Arc<Notify>,
}
pub(crate) fn channel() -> (CapturePublisher, CaptureWatch) {
let (tx, rx) = watch::channel(None);
let nudge = Arc::new(Notify::new());
(
CapturePublisher {
tx: Arc::new(tx),
nudge: nudge.clone(),
},
CaptureWatch { rx, nudge },
)
}
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub(crate) struct HostDelivery {
pub(crate) delivered: bool,
pub(crate) since: Option<i64>,
pub(crate) reason: Option<String>,
pub(crate) core: Option<String>,
pub(crate) customer_override: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub(crate) struct DeliveryState {
pub(crate) hub: Option<HostDelivery>,
pub(crate) vscode: Option<HostDelivery>,
pub(crate) hub_desired: Option<Desired>,
}
pub(crate) fn state_path(ol_dir: &Path) -> PathBuf {
ol_dir.join("state").join(STATE_FILE)
}
pub(crate) fn read_state(ol_dir: &Path) -> Option<DeliveryState> {
let raw = std::fs::read_to_string(state_path(ol_dir)).ok()?;
serde_json::from_str(&raw).ok()
}
pub(crate) fn hub_delivered(ol_dir: &Path) -> bool {
read_state(ol_dir)
.and_then(|state| state.hub)
.is_some_and(|hub| hub.delivered)
}
pub(crate) fn read_hub_desired(ol_dir: &Path) -> Desired {
read_state(ol_dir)
.and_then(|s| s.hub_desired)
.unwrap_or_default()
}
fn write_state(ol_dir: &Path, state: &DeliveryState) -> bool {
match crate::fs_secure::write_json_state(&state_path(ol_dir), state) {
Ok(()) => true,
Err(e) => {
tracing::warn!(
code = crate::error::ERR_CLINE_PLUGIN_NOT_LOADED,
error = %e,
"could not write the Cline plugin delivery state"
);
false
}
}
}
#[derive(Debug, Clone)]
pub(crate) struct TickInputs {
pub(crate) binding: bool,
pub(crate) owns: bool,
pub(crate) hub_managed: bool,
pub(crate) app: Option<Option<ClineApp>>,
pub(crate) assessments: Vec<HostAssessment>,
pub(crate) capture: CaptureState,
}
pub(crate) struct HubSeam<'a> {
pub(crate) table: &'a dyn ProcessTable,
pub(crate) port: u16,
pub(crate) recheck: Duration,
pub(crate) poll: Duration,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct TickReport {
pub(crate) state: DeliveryState,
pub(crate) hub_action: Option<HubAction>,
pub(crate) hub_outcome: Option<HubOutcome>,
pub(crate) login_env: Option<LoginEnvOutcome>,
}
#[cfg(test)]
pub(crate) fn tick_with(
sink: &dyn Sink,
hub: &HubSeam<'_>,
ol_dir: &Path,
inputs: TickInputs,
now_ms: i64,
) -> TickReport {
tick_from(sink, hub, ol_dir, inputs, now_ms, &mut read_state(ol_dir))
}
fn tick_from(
sink: &dyn Sink,
hub: &HubSeam<'_>,
ol_dir: &Path,
inputs: TickInputs,
now_ms: i64,
previous: &mut Option<DeliveryState>,
) -> TickReport {
if !inputs.binding {
let login_env = inputs.owns.then(|| {
let mut ledger = Ledger::load(ol_dir);
login_env::reconcile_login_env(sink, &mut ledger, None)
});
let hub_outcome = inputs.hub_managed.then(|| {
cline_hub::release_with(
hub.table,
inputs.app.as_ref().and_then(Option::as_ref),
hub.port,
hub.recheck,
hub.poll,
)
});
return TickReport {
state: DeliveryState::default(),
hub_action: None,
hub_outcome,
login_env,
};
}
let empty = DeliveryState::default();
let before = previous.as_ref().unwrap_or(&empty);
let find = |class: HostClass| inputs.assessments.iter().find(|a| a.host.class == class);
let mut hub_action = None;
let mut hub_outcome = None;
let mut hub_desired = None;
let hub_entry = match (&inputs.app, find(HostClass::Hub)) {
(Some(None), _) => Some(HostDelivery {
reason: Some(reason::NOT_DETECTED_ON_PLATFORM.to_string()),
..HostDelivery::default()
}),
(_, None) => None,
(_, Some(a)) if !inputs.hub_managed => Some(entry(a, Some(reason::NOT_MANAGED.into()))),
(_, Some(a)) if inputs.capture.is_none() => {
hub_desired = before.hub_desired.clone();
Some(before.hub.clone().unwrap_or_else(|| {
entry(
a,
Some(format!(
"{}: waiting for the model relay's first report",
reason::PENDING_RESTART
)),
)
}))
}
(app, Some(a)) => {
let desired = Desired {
capture: inputs.capture.clone().flatten(),
plugin: a
.deliverable()
.then(|| a.wrapper.clone())
.flatten()
.map(|runtime| PluginEnv {
wrapper_path: crate::hooks::cline_bootstrap::wrapper_path(ol_dir),
runtime,
}),
};
let outcome = match app {
Some(Some(app)) => {
let (action, outcome) = cline_hub::reconcile(
hub.table,
app,
hub.port,
&desired,
hub.recheck,
hub.poll,
);
hub_action = Some(action);
Some(outcome)
}
_ => None,
};
let why = a.reason.clone().or_else(|| match &outcome {
Some(o) if o.carries_desired() => None,
Some(HubOutcome::DeferredBusy | HubOutcome::IdleUnknown) => {
Some(reason::PENDING_RESTART.to_string())
}
Some(HubOutcome::Failed(detail)) => {
Some(format!("{}: hub: {detail}", reason::PROBE_FAILED))
}
Some(_) => Some(format!(
"{}: hub: the port is held by a process OpenLatch does not manage",
reason::NOT_MANAGED
)),
None => Some(format!(
"{}: hub: Cline.app not found",
reason::PROBE_FAILED
)),
});
hub_outcome = outcome;
hub_desired = Some(desired);
Some(entry(a, why))
}
};
let vscode = find(HostClass::VsCode);
let (vscode_entry, login_env) = if !inputs.owns {
(
vscode.map(|a| entry(a, Some(reason::NOT_MANAGED.into()))),
None,
)
} else {
let wrapper_path = crate::hooks::cline_bootstrap::wrapper_path(ol_dir);
let want = vscode
.filter(|a| a.deliverable())
.and_then(|a| a.wrapper.as_ref())
.and_then(|runtime| Some((wrapper_path.to_str()?, runtime.to_str()?)));
let mut ledger = Ledger::load(ol_dir);
let outcome = login_env::reconcile_login_env(sink, &mut ledger, want);
let entry = vscode.map(|a| {
let mut e = entry(
a,
match &outcome {
LoginEnvOutcome::Delivered | LoginEnvOutcome::Released => a.reason.clone(),
LoginEnvOutcome::CustomerOverride { names } => Some(format!(
"{}: {}",
reason::CUSTOMER_OVERRIDE,
names.join(", ")
)),
LoginEnvOutcome::Failed { detail } => {
Some(format!("{}: login-env: {detail}", reason::PROBE_FAILED))
}
},
);
e.customer_override = matches!(outcome, LoginEnvOutcome::CustomerOverride { .. });
if e.reason.is_none() && outcome != LoginEnvOutcome::Delivered {
e.reason = Some(format!("{}: login-env", reason::PROBE_FAILED));
}
e
});
(entry, Some(outcome))
};
let stamp = |e: Option<HostDelivery>, before: &Option<HostDelivery>| {
e.map(|mut e| {
e.delivered = e.reason.is_none();
e.since = if e.delivered {
before
.as_ref()
.filter(|b| b.delivered)
.and_then(|b| b.since)
.or(Some(now_ms))
} else {
None
};
e
})
};
let state = DeliveryState {
hub: stamp(hub_entry, &before.hub),
vscode: stamp(vscode_entry, &before.vscode),
hub_desired,
};
if previous.as_ref() != Some(&state) && write_state(ol_dir, &state) {
*previous = Some(state.clone());
}
TickReport {
state,
hub_action,
hub_outcome,
login_env,
}
}
fn entry(a: &HostAssessment, reason: Option<String>) -> HostDelivery {
HostDelivery {
delivered: false,
since: None,
reason,
core: a.host.core_version.as_ref().map(ToString::to_string),
customer_override: false,
}
}
fn gather(config: &Config, ol_dir: &Path, capture: CaptureState) -> TickInputs {
let owns_machine = crate::supervision::owns_machine_supervision();
let cline = crate::hooks::detect_agents()
.into_iter()
.find(|a| a.agent_type() == ClineBinding::AGENT_TYPE);
match cline {
None => {
let hub_managed = cline_hub::manages_hub_without_binding(config);
TickInputs {
binding: false,
owns: config.model_relay.owns_agent_wiring() && owns_machine,
hub_managed,
app: hub_managed.then(cline_app::detect),
assessments: Vec::new(),
capture,
}
}
Some(a) => {
let detected = cline_app::detect();
TickInputs {
binding: true,
owns: super::owns_wiring_for(config, &*a.binding) && owns_machine,
hub_managed: cline_hub::manages_hub(config, &*a.binding),
app: cline_app::visible_app_from(detected.clone()),
assessments: cline_runtime::assess_hosts_with(
ol_dir,
config.cline.plugin_delivery,
detected,
),
capture,
}
}
}
}
static TICK_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
static STOPPED: AtomicBool = AtomicBool::new(false);
fn tick(
config: &Config,
ol_dir: &Path,
capture: CaptureState,
sink: &dyn Sink,
previous: &mut Option<DeliveryState>,
) {
let _held = TICK_LOCK.lock().unwrap_or_else(|e| e.into_inner());
if STOPPED.load(Ordering::SeqCst) {
return;
}
let inputs = gather(config, ol_dir, capture);
let table = cline_app::process::table();
let seam = HubSeam {
table: &*table,
port: cline_app::hub_port(),
recheck: cline_app::IDLE_RECHECK,
poll: cline_app::HUB_STARTUP_POLL,
};
let report = tick_from(sink, &seam, ol_dir, inputs, now_ms(), previous);
tracing::debug!(
hub = ?report.state.hub,
vscode = ?report.state.vscode,
hub_action = ?report.hub_action,
hub_outcome = ?report.hub_outcome,
login_env = ?report.login_env,
"Cline plugin delivery tick"
);
}
fn now_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| i64::try_from(d.as_millis()).unwrap_or(i64::MAX))
}
pub(crate) fn release_on_stop(config: &Config) {
STOPPED.store(true, Ordering::SeqCst);
let _held = TICK_LOCK.lock().unwrap_or_else(|e| e.into_inner());
cline_hub::release_hub(config);
}
async fn run(
config: Arc<Config>,
mut capture: watch::Receiver<CaptureState>,
nudge: Arc<Notify>,
sink: &'static dyn Sink,
mut shutdown: watch::Receiver<bool>,
) {
let ol_dir = crate::config::openlatch_dir();
let previous = Arc::new(std::sync::Mutex::new(read_state(&ol_dir)));
loop {
if *shutdown.borrow() {
return;
}
let state = capture.borrow_and_update().clone();
let cfg = config.clone();
let dir = ol_dir.clone();
let prev = previous.clone();
let ticked = tokio::task::spawn_blocking(move || {
let mut prev = prev.lock().unwrap_or_else(|e| e.into_inner());
tick(&cfg, &dir, state, sink, &mut prev);
});
if let Err(e) = ticked.await {
tracing::warn!(error = %e, "Cline plugin delivery tick did not complete");
}
tokio::select! {
_ = shutdown.wait_for(|stop| *stop) => return,
_ = tokio::time::sleep(TICK) => {}
_ = nudge.notified() => {}
}
}
}
pub(crate) fn spawn(
health: &Arc<HealthRegistry>,
config: Arc<Config>,
watch: CaptureWatch,
sink: &'static dyn Sink,
shutdown: watch::Receiver<bool>,
) -> tokio::task::JoinHandle<()> {
let task_shutdown = shutdown.clone();
spawn_supervised(
health,
TaskSpec::new(
super::subsystem::CLINE_PLUGIN_DELIVERY,
RestartPolicy::Always,
),
shutdown,
move || {
run(
config.clone(),
watch.rx.clone(),
watch.nudge.clone(),
sink,
task_shutdown.clone(),
)
},
)
}
#[cfg(test)]
mod tests {
use std::path::{Path, PathBuf};
use std::time::Duration;
use super::*;
use crate::core::login_env::test_support::{Call, RecordingSink};
use crate::core::login_env::ALLOWED;
use crate::hooks::cline_app::process::test_support::FakeTable;
use crate::hooks::cline_app::test_fixture::{
clean_env, hub_fx, listen, record_exists, seed_record_for, HubFx,
};
use crate::hooks::cline_hosts::Host;
use crate::hooks::cline_runtime::ProbeResult;
const RELAY: &str = "http://127.0.0.1:7600";
fn host(class: HostClass, runtime: PathBuf) -> Host {
Host {
class,
runtime: runtime.clone(),
runtime_var: if class == HostClass::Hub {
"BUN_BE_BUN"
} else {
"ELECTRON_RUN_AS_NODE"
},
core_version: Some(semver::Version::new(0, 0, 86)),
source: runtime,
}
}
fn deliverable(ol: &Path, class: HostClass, runtime: PathBuf) -> HostAssessment {
HostAssessment {
host: host(class, runtime),
wrapper: Some(cline_runtime::wrapper_for(ol, class)),
reason: None,
}
}
fn with_reason(mut a: HostAssessment, reason: &str) -> HostAssessment {
a.reason = Some(reason.to_string());
a
}
fn both_deliverable(fx: &HubFx, ol: &Path) -> Vec<HostAssessment> {
vec![
deliverable(ol, HostClass::Hub, fx.app.sidecar.clone()),
deliverable(ol, HostClass::VsCode, fx.root.path().join("Code")),
]
}
fn inputs(fx: &HubFx, assessments: Vec<HostAssessment>, capture: CaptureState) -> TickInputs {
TickInputs {
binding: true,
owns: true,
hub_managed: true,
app: Some(Some(fx.app.clone())),
assessments,
capture,
}
}
fn run(
sink: &RecordingSink,
t: &FakeTable,
fx: &HubFx,
ol: &Path,
i: TickInputs,
) -> TickReport {
let seam = HubSeam {
table: t,
port: fx.port,
recheck: Duration::ZERO,
poll: Duration::ZERO,
};
tick_with(sink, &seam, ol, i, 1_000)
}
fn plugin(ol: &Path) -> PluginEnv {
PluginEnv {
wrapper_path: crate::hooks::cline_bootstrap::wrapper_path(ol),
runtime: cline_runtime::wrapper_for(ol, HostClass::Hub),
}
}
fn starts_as(fx: &HubFx, t: &FakeTable, pid: u32) {
*t.on_spawn_listen.lock().unwrap() = Some(pid);
*t.next_pid.lock().unwrap() = pid;
t.procs
.lock()
.unwrap()
.insert(pid, fx.hub_info(Some(clean_env())));
}
fn login_values(sink: &RecordingSink) -> [Option<String>; 2] {
[sink.value(ALLOWED[0]), sink.value(ALLOWED[1])]
}
fn code(reason: Option<&String>) -> Option<&str> {
reason.map(|r| r.split(':').next().unwrap_or(r))
}
#[test]
fn unpublished_capture_never_restarts_the_hub() {
let fx = hub_fx();
let ol = crate::config::openlatch_dir();
let t = fx.fake();
seed_record_for(
&fx,
500,
1,
"0.0.34",
&crate::hooks::cline_app::test_fixture::captured(RELAY),
);
listen(&t, 500, fx.hub_info(Some(clean_env())));
*t.clients.lock().unwrap() = Ok(vec![500, 4242]);
let sink = RecordingSink::default();
let report = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, both_deliverable(&fx, &ol), None),
);
assert_eq!(
report.hub_action, None,
"no hub action on an unreported tick"
);
assert_eq!(report.hub_outcome, None);
assert!(t.terminated.lock().unwrap().is_empty());
assert!(t.spawns.lock().unwrap().is_empty());
assert!(record_exists(), "the record is untouched");
assert_eq!(report.login_env, Some(LoginEnvOutcome::Delivered));
let wrapper = crate::hooks::cline_bootstrap::wrapper_path(&ol);
let runtime = cline_runtime::wrapper_for(&ol, HostClass::VsCode);
assert_eq!(
login_values(&sink),
[
Some(wrapper.display().to_string()),
Some(runtime.display().to_string())
]
);
let hub = report.state.hub.as_ref().expect("a hub row");
assert!(!hub.delivered);
assert_eq!(code(hub.reason.as_ref()), Some(reason::PENDING_RESTART));
assert!(
report
.state
.vscode
.as_ref()
.expect("a vscode row")
.delivered
);
assert_eq!(read_state(&ol), Some(report.state));
}
#[test]
fn relay_off_takes_capture_off_now_keeping_plugin() {
let fx = hub_fx();
let ol = crate::config::openlatch_dir();
let t = fx.fake();
seed_record_for(
&fx,
500,
1,
"0.0.34",
&Desired {
capture: crate::hooks::cline_app::test_fixture::captured(RELAY).capture,
plugin: Some(plugin(&ol)),
},
);
listen(&t, 500, fx.hub_info(Some(clean_env())));
*t.clients.lock().unwrap() = Ok(vec![500, 4242]);
starts_as(&fx, &t, 900);
let sink = RecordingSink::default();
let report = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, both_deliverable(&fx, &ol), Some(None)),
);
assert_eq!(report.hub_action, Some(HubAction::RestartNow));
assert_eq!(report.hub_outcome, Some(HubOutcome::Restarted));
assert_eq!(*t.terminated.lock().unwrap(), vec![500]);
let spawns = t.spawns.lock().unwrap().clone();
assert_eq!(spawns.len(), 1);
let want = Desired {
capture: None,
plugin: Some(plugin(&ol)),
};
assert_eq!(spawns[0].3, cline_app::hub_env(fx.port, &want));
let names: Vec<&str> = spawns[0].3.iter().map(|(k, _)| k.as_str()).collect();
assert!(ALLOWED.iter().all(|n| names.contains(n)), "{names:?}");
assert!(
!names.iter().any(|n| n.eq_ignore_ascii_case("HTTPS_PROXY")),
"{names:?}"
);
assert_eq!(report.state.hub_desired, Some(want));
let hub = report.state.hub.as_ref().expect("a hub row");
assert!(hub.delivered);
assert_eq!(hub.since, Some(1_000));
}
#[test]
fn delivery_off_clears_ours_and_restarts_hub_if_idle() {
let fx = hub_fx();
let ol = crate::config::openlatch_dir();
let t = fx.fake();
let capture = crate::hooks::cline_app::ProxyEnvVars::for_relay(7600);
let sink = RecordingSink::default();
starts_as(&fx, &t, 900);
let first = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, both_deliverable(&fx, &ol), Some(Some(capture.clone()))),
);
assert_eq!(first.hub_action, Some(HubAction::Start));
assert_eq!(first.hub_outcome, Some(HubOutcome::Started));
assert_eq!(first.login_env, Some(LoginEnvOutcome::Delivered));
assert!(login_values(&sink).iter().all(Option::is_some));
*t.clients.lock().unwrap() = Ok(vec![900, 4242]);
let off = || {
both_deliverable(&fx, &ol)
.into_iter()
.map(|a| with_reason(a, reason::DELIVERY_OFF))
.collect::<Vec<_>>()
};
let busy = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, off(), Some(Some(capture.clone()))),
);
assert_eq!(busy.hub_action, Some(HubAction::RestartIfIdle));
assert_eq!(busy.hub_outcome, Some(HubOutcome::DeferredBusy));
assert!(t.terminated.lock().unwrap().is_empty(), "never Now");
assert_eq!(busy.login_env, Some(LoginEnvOutcome::Released));
assert_eq!(
login_values(&sink),
[None, None],
"our login values are cleared"
);
assert!(sink
.calls()
.iter()
.any(|c| matches!(c, Call::Clear(n) if n == ALLOWED[0])));
*t.clients.lock().unwrap() = Ok(Vec::new());
starts_as(&fx, &t, 901);
let idle = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, off(), Some(Some(capture.clone()))),
);
assert_eq!(idle.hub_action, Some(HubAction::RestartIfIdle));
assert_eq!(idle.hub_outcome, Some(HubOutcome::Restarted));
assert_eq!(*t.terminated.lock().unwrap(), vec![900]);
let spawns = t.spawns.lock().unwrap().clone();
assert_eq!(spawns.len(), 2);
let want = Desired {
capture: Some(capture.clone()),
plugin: None,
};
assert_eq!(
spawns[1].3,
cline_app::hub_env(fx.port, &want),
"capture untouched"
);
assert_eq!(idle.state.hub_desired, Some(want));
for row in [idle.state.hub.as_ref(), idle.state.vscode.as_ref()] {
let row = row.expect("a row");
assert!(!row.delivered);
assert_eq!(row.reason.as_deref(), Some(reason::DELIVERY_OFF));
}
assert!(Ledger::load(&ol).is_empty());
}
#[test]
fn delivery_tick_removes_values_when_runtime_disappears() {
let fx = hub_fx();
let ol = crate::config::openlatch_dir();
let t = fx.fake();
let sink = RecordingSink::default();
let code_bin = fx.root.path().join("Code");
let hub_bin = fx.root.path().join("hub-runtime");
std::fs::write(&code_bin, "").expect("a fake VS Code runtime");
std::fs::write(&hub_bin, "").expect("a fake hub runtime");
starts_as(&fx, &t, 900);
let first = run(
&sink,
&t,
&fx,
&ol,
inputs(
&fx,
vec![
deliverable(&ol, HostClass::Hub, hub_bin.clone()),
deliverable(&ol, HostClass::VsCode, code_bin.clone()),
],
Some(None),
),
);
assert_eq!(first.hub_action, Some(HubAction::Start));
assert_eq!(first.login_env, Some(LoginEnvOutcome::Delivered));
assert!(first.state.hub.as_ref().unwrap().delivered);
assert!(first.state.vscode.as_ref().unwrap().delivered);
std::fs::remove_file(&code_bin).expect("remove the VS Code runtime");
std::fs::remove_file(&hub_bin).expect("remove the hub runtime");
let probe_ran = std::sync::atomic::AtomicBool::new(false);
let probe = |_: &Path, _: &Host| {
probe_ran.store(true, Ordering::SeqCst);
ProbeResult {
ok: true,
stage: None,
error: None,
}
};
let assessments: Vec<HostAssessment> = [
host(HostClass::Hub, hub_bin.clone()),
host(HostClass::VsCode, code_bin.clone()),
]
.iter()
.map(|h| cline_runtime::assess_with(&ol, h, true, &Ok(()), &probe))
.collect();
assert!(
!probe_ran.load(Ordering::SeqCst),
"no probe past a missing runtime"
);
for a in &assessments {
assert_eq!(a.reason.as_deref(), Some(reason::RUNTIME_MISSING), "{a:?}");
}
let second = run(&sink, &t, &fx, &ol, inputs(&fx, assessments, Some(None)));
assert_eq!(second.login_env, Some(LoginEnvOutcome::Released));
assert_eq!(login_values(&sink), [None, None]);
assert!(Ledger::load(&ol).is_empty());
assert_eq!(second.hub_action, Some(HubAction::RestartIfIdle));
assert_eq!(second.hub_outcome, Some(HubOutcome::Restarted));
assert_eq!(*t.terminated.lock().unwrap(), vec![900]);
assert_eq!(t.spawns.lock().unwrap().len(), 1, "only the first start");
assert!(!record_exists());
assert_eq!(second.state.hub_desired, Some(Desired::default()));
for row in [second.state.hub.as_ref(), second.state.vscode.as_ref()] {
let row = row.expect("a row");
assert!(!row.delivered);
assert_eq!(row.since, None);
assert_eq!(row.reason.as_deref(), Some(reason::RUNTIME_MISSING));
}
}
#[test]
fn hub_is_delivered_only_once_it_carries_the_env() {
let fx = hub_fx();
let ol = crate::config::openlatch_dir();
let t = fx.fake();
listen(&t, 500, fx.hub_info(Some(clean_env())));
*t.clients.lock().unwrap() = Ok(vec![500, 4242]);
let sink = RecordingSink::default();
let busy = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, both_deliverable(&fx, &ol), Some(None)),
);
assert_eq!(busy.hub_action, Some(HubAction::RestartIfIdle));
assert_eq!(busy.hub_outcome, Some(HubOutcome::DeferredBusy));
let hub = busy.state.hub.as_ref().expect("a hub row");
assert!(!hub.delivered, "a deferred restart is not a delivery");
assert_eq!(hub.since, None);
assert_eq!(hub.reason.as_deref(), Some(reason::PENDING_RESTART));
assert_eq!(read_state(&ol), Some(busy.state.clone()));
*t.clients.lock().unwrap() = Ok(Vec::new());
starts_as(&fx, &t, 900);
let idle = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, both_deliverable(&fx, &ol), Some(None)),
);
assert_eq!(idle.hub_outcome, Some(HubOutcome::Restarted));
let hub = idle.state.hub.as_ref().expect("a hub row");
assert!(hub.delivered);
assert_eq!(hub.reason, None);
assert_eq!(hub.since, Some(1_000));
}
#[test]
fn binding_gone_releases_the_hub_this_tick() {
let fx = hub_fx();
let ol = crate::config::openlatch_dir();
let t = fx.fake();
let sink = RecordingSink::default();
starts_as(&fx, &t, 900);
let first = run(
&sink,
&t,
&fx,
&ol,
inputs(&fx, both_deliverable(&fx, &ol), Some(None)),
);
assert_eq!(first.hub_outcome, Some(HubOutcome::Started));
assert!(login_values(&sink).iter().all(Option::is_some));
*t.clients.lock().unwrap() = Ok(vec![900, 4242]);
let mut gone = inputs(&fx, Vec::new(), Some(None));
gone.binding = false;
let report = run(&sink, &t, &fx, &ol, gone);
assert_eq!(report.login_env, Some(LoginEnvOutcome::Released));
assert_eq!(login_values(&sink), [None, None]);
assert_eq!(report.hub_outcome, Some(HubOutcome::Restarted));
assert_eq!(*t.terminated.lock().unwrap(), vec![900]);
let spawns = t.spawns.lock().unwrap().clone();
assert_eq!(spawns.len(), 2, "the first start, then the clean respawn");
assert_eq!(
spawns[1].3,
cline_app::hub_env(fx.port, &Desired::default()),
"released with the empty Desired"
);
assert!(!record_exists());
let t = fx.fake();
let mut gone = inputs(&fx, Vec::new(), Some(None));
gone.binding = false;
gone.hub_managed = false;
let report = run(&sink, &t, &fx, &ol, gone);
assert_eq!(report.hub_outcome, None);
assert!(t.spawns.lock().unwrap().is_empty() && t.terminated.lock().unwrap().is_empty());
}
#[test]
fn isolated_instance_never_touches_login_env() {
let fx = hub_fx();
let ol = crate::config::openlatch_dir();
assert!(!crate::supervision::owns_machine_supervision());
let sink =
RecordingSink::with(&[(ALLOWED[0], "/old/wrapper"), (ALLOWED[1], "/old/runtime")]);
let t = fx.fake();
for (binding, capture) in [(true, Some(None)), (true, None), (false, Some(None))] {
let mut i = inputs(&fx, both_deliverable(&fx, &ol), capture);
i.binding = binding;
i.owns = false;
i.hub_managed = false;
let report = run(&sink, &t, &fx, &ol, i);
assert_eq!(report.login_env, None, "binding={binding}");
if binding {
let vscode = report.state.vscode.as_ref().expect("a vscode row");
assert!(!vscode.delivered);
assert_eq!(code(vscode.reason.as_ref()), Some(reason::NOT_MANAGED));
}
}
assert_eq!(sink.calls(), Vec::<Call>::new(), "zero calls on the sink");
assert!(t.spawns.lock().unwrap().is_empty());
assert!(!Ledger::path_in(&ol).exists());
}
}