fix(apps): preserve lifecycle state and wait for usable launch endpoints
This commit is contained in:
@@ -114,6 +114,9 @@ impl PortMap {
|
||||
/// there.
|
||||
fn apps_dirs() -> Vec<PathBuf> {
|
||||
let mut dirs = Vec::new();
|
||||
if let Some(root) = std::env::var_os("ARCHIPELAGO_APPS_DIR") {
|
||||
dirs.push(root.into());
|
||||
}
|
||||
if let Ok(manifest_dir) = std::env::var("CARGO_MANIFEST_DIR") {
|
||||
dirs.push(PathBuf::from(manifest_dir).join("../../apps"));
|
||||
}
|
||||
|
||||
@@ -144,6 +144,34 @@ pub fn shared_status() -> Arc<RwLock<GateStatus>> {
|
||||
.clone()
|
||||
}
|
||||
|
||||
static REFRESH_KICK: std::sync::LazyLock<tokio::sync::Notify> =
|
||||
std::sync::LazyLock::new(tokio::sync::Notify::new);
|
||||
static REFRESH_REV: std::sync::LazyLock<tokio::sync::watch::Sender<u64>> =
|
||||
std::sync::LazyLock::new(|| tokio::sync::watch::channel(0).0);
|
||||
|
||||
/// Installation must not wait for the minute sweep before becoming reachable.
|
||||
/// Wait for a completed sweep, bounded if shutdown/startup prevents one.
|
||||
pub async fn refresh_now() {
|
||||
let mut completed = REFRESH_REV.subscribe();
|
||||
REFRESH_KICK.notify_one();
|
||||
let _ = tokio::time::timeout(std::time::Duration::from_secs(3), completed.changed()).await;
|
||||
}
|
||||
|
||||
pub fn port_claimed(status: &GateStatus, port: u16) -> bool {
|
||||
let mut external = false;
|
||||
let mut tor = false;
|
||||
for (claimed_port, address) in &status.claimed {
|
||||
if *claimed_port != port {
|
||||
continue;
|
||||
}
|
||||
if let Ok(ip) = address.parse::<IpAddr>() {
|
||||
tor |= ip == GATE_TOR_UPSTREAM;
|
||||
external |= !ip.is_loopback();
|
||||
}
|
||||
}
|
||||
external && tor
|
||||
}
|
||||
|
||||
/// Run the gate. Returns only on shutdown.
|
||||
pub async fn run(
|
||||
gate: Arc<AppGate>,
|
||||
@@ -162,11 +190,12 @@ pub async fn run(
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = interval.tick() => {
|
||||
sweep(&gate, &status, &mut held, &shutdown_rx).await;
|
||||
}
|
||||
_ = interval.tick() => {}
|
||||
_ = REFRESH_KICK.notified() => {}
|
||||
_ = shutdown_rx.changed() => return,
|
||||
}
|
||||
sweep(&gate, &status, &mut held, &shutdown_rx).await;
|
||||
REFRESH_REV.send_modify(|revision| *revision = revision.wrapping_add(1));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -461,3 +490,19 @@ mod tests {
|
||||
assert!(!status.is_fully_enforced());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod readiness_tests {
|
||||
use super::*;
|
||||
#[test]
|
||||
fn readiness_requires_external_and_tor_claims_for_the_same_port() {
|
||||
let mut status = GateStatus::default();
|
||||
assert!(!port_claimed(&status, 3001));
|
||||
status.claimed.push((3001, "127.0.0.2".into()));
|
||||
assert!(!port_claimed(&status, 3001));
|
||||
status.claimed.push((3002, "192.0.2.10".into()));
|
||||
assert!(!port_claimed(&status, 3001));
|
||||
status.claimed.push((3001, "192.0.2.10".into()));
|
||||
assert!(port_claimed(&status, 3001));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user