use crate::api::AppState;
use crate::error::Result;
use crate::guest_query::list_running_container_ids;
use crate::handlers::setup_container_networking;
use crate::proxy::ProxyState;
use arcbox_core::{Runtime, VmLifecycleState};
use std::collections::HashSet;
use std::future::Future;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::watch;
use tokio_util::sync::CancellationToken;
const RECONCILE_INTERVAL: Duration = Duration::from_secs(30);
const RETURN_RETRY: Duration = Duration::from_secs(2);
pub fn spawn(runtime: Arc<Runtime>, proxy: Arc<ProxyState>, shutdown: CancellationToken) {
let vm = runtime.subscribe_system_vm_state();
let guest = ProxyGuest {
state: AppState {
runtime: Arc::clone(&runtime),
proxy,
},
};
drop(tokio::spawn(async move {
reconcile_loop(&runtime, &guest, vm, shutdown).await;
}));
}
trait Guest: Send + Sync {
fn running(&self) -> impl Future<Output = Result<HashSet<String>>> + Send;
fn set_up(&self, container_id: &str) -> impl Future<Output = ()> + Send;
}
struct ProxyGuest {
state: AppState,
}
impl Guest for ProxyGuest {
async fn running(&self) -> Result<HashSet<String>> {
self.state
.proxy
.reset_if_restarted(self.state.runtime.system_vm_restart_generation());
list_running_container_ids(self.state.proxy.client()).await
}
async fn set_up(&self, container_id: &str) {
setup_container_networking(&self.state, container_id).await;
}
}
async fn reconcile_loop(
runtime: &Runtime,
guest: &impl Guest,
mut vm: watch::Receiver<VmLifecycleState>,
shutdown: CancellationToken,
) {
let mut ready = vm.borrow_and_update().is_ready();
let mut returning = true;
loop {
let delay = if returning {
RETURN_RETRY
} else {
RECONCILE_INTERVAL
};
tokio::select! {
() = shutdown.cancelled() => break,
changed = vm.changed() => {
if changed.is_err() {
break;
}
let now_ready = vm.borrow_and_update().is_ready();
if ready && !now_ready {
runtime.retire_system_vm_container_networking().await;
}
if now_ready && !ready {
returning = true;
}
ready = now_ready;
}
() = tokio::time::sleep(delay), if ready => {
match reconcile(runtime, guest).await {
Ok(()) => returning = false,
Err(e) => tracing::debug!(error = %e, "host networking reconcile skipped"),
}
}
}
}
}
async fn reconcile(runtime: &Runtime, guest: &impl Guest) -> Result<()> {
let registered = runtime.registered_container_ids().await;
let running = guest.running().await?;
for id in registered.difference(&running) {
tracing::info!(
container_id = %id,
"reconciler tearing down host networking for a container no longer running"
);
runtime.stop_port_forwarding_by_id(id).await;
runtime.deregister_dns_by_id(id).await;
}
for id in running.difference(®istered) {
tracing::info!(
container_id = %id,
"reconciler setting up host networking for a running container it did not know"
);
guest.set_up(id).await;
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::error::DockerError;
use arcbox_core::{Config, Runtime, VmLifecycleConfig};
use std::net::{IpAddr, Ipv4Addr};
use std::sync::Mutex;
use tempfile::TempDir;
fn test_runtime() -> (Arc<Runtime>, TempDir) {
let tmp = TempDir::new().unwrap();
let config = Config {
data_dir: tmp.path().to_path_buf(),
..Default::default()
};
let vlc = VmLifecycleConfig {
skip_vm_check: true,
..Default::default()
};
let runtime = Arc::new(Runtime::with_vm_lifecycle_config(config, vlc).expect("runtime"));
(runtime, tmp)
}
struct FakeGuest {
running: std::result::Result<HashSet<String>, String>,
set_up: Mutex<Vec<String>>,
}
impl FakeGuest {
fn running(ids: &[&str]) -> Self {
Self {
running: Ok(ids.iter().map(|id| (*id).to_owned()).collect()),
set_up: Mutex::new(Vec::new()),
}
}
fn unreachable() -> Self {
Self {
running: Err("guest unreachable".to_owned()),
set_up: Mutex::new(Vec::new()),
}
}
fn set_up(&self) -> Vec<String> {
self.set_up.lock().unwrap().clone()
}
}
impl Guest for FakeGuest {
async fn running(&self) -> Result<HashSet<String>> {
self.running.clone().map_err(DockerError::Server)
}
async fn set_up(&self, container_id: &str) {
self.set_up.lock().unwrap().push(container_id.to_owned());
}
}
#[tokio::test]
async fn reconcile_tears_down_only_orphans() {
let (runtime, _tmp) = test_runtime();
let ip = IpAddr::V4(Ipv4Addr::LOCALHOST);
runtime
.register_dns("alive", &["alive.local".into()], ip)
.await;
runtime
.register_dns("dead", &["dead.local".into()], ip)
.await;
assert_eq!(runtime.registered_container_ids().await.len(), 2);
let guest = FakeGuest::running(&["alive"]);
reconcile(&runtime, &guest).await.unwrap();
let remaining = runtime.registered_container_ids().await;
assert!(remaining.contains("alive"));
assert!(!remaining.contains("dead"));
assert_eq!(guest.set_up(), Vec::<String>::new());
}
#[tokio::test]
async fn reconcile_reclaims_alias_only_containers() {
let (runtime, _tmp) = test_runtime();
runtime
.register_container_alias("ephemeral", "cafe1234")
.await;
assert!(
runtime
.registered_container_ids()
.await
.contains("cafe1234"),
"alias-only containers must be visible to the reconciler"
);
reconcile(&runtime, &FakeGuest::running(&[])).await.unwrap();
assert!(runtime.registered_container_ids().await.is_empty());
assert_eq!(
runtime.resolve_registered_container("ephemeral").await,
None
);
}
#[tokio::test]
async fn reconcile_is_fail_safe_on_query_error() {
let (runtime, _tmp) = test_runtime();
runtime
.register_dns("c", &["c.local".into()], IpAddr::V4(Ipv4Addr::LOCALHOST))
.await;
let guest = FakeGuest::unreachable();
let result = reconcile(&runtime, &guest).await;
assert!(result.is_err());
assert!(runtime.registered_container_ids().await.contains("c"));
assert_eq!(guest.set_up(), Vec::<String>::new());
}
#[tokio::test]
async fn reconcile_sets_up_running_containers_the_host_does_not_know() {
let (runtime, _tmp) = test_runtime();
runtime
.register_dns(
"known",
&["known.local".into()],
IpAddr::V4(Ipv4Addr::LOCALHOST),
)
.await;
let guest = FakeGuest::running(&["known", "restarted"]);
reconcile(&runtime, &guest).await.unwrap();
assert_eq!(guest.set_up(), ["restarted"]);
assert!(runtime.registered_container_ids().await.contains("known"));
}
async fn settle(elapsed: Duration) {
tokio::time::sleep(elapsed).await;
tokio::task::yield_now().await;
}
#[tokio::test(start_paused = true)]
async fn a_vm_departure_retires_container_state_and_its_return_is_reconciled_promptly() {
let (runtime, _tmp) = test_runtime();
let guest = Arc::new(FakeGuest::running(&["web"]));
runtime
.register_dns(
"web",
&["web.local".into()],
IpAddr::V4(Ipv4Addr::LOCALHOST),
)
.await;
let (vm_tx, vm) = watch::channel(VmLifecycleState::Running);
let shutdown = CancellationToken::new();
let task = tokio::spawn({
let (runtime, guest, shutdown) =
(Arc::clone(&runtime), Arc::clone(&guest), shutdown.clone());
async move { reconcile_loop(&runtime, guest.as_ref(), vm, shutdown).await }
});
settle(RECONCILE_INTERVAL * 2).await;
assert_eq!(guest.set_up(), Vec::<String>::new());
assert!(runtime.registered_container_ids().await.contains("web"));
vm_tx.send_replace(VmLifecycleState::Stopping);
settle(Duration::ZERO).await;
assert!(runtime.registered_container_ids().await.is_empty());
settle(RECONCILE_INTERVAL * 2).await;
assert_eq!(guest.set_up(), Vec::<String>::new());
vm_tx.send_replace(VmLifecycleState::Running);
settle(RETURN_RETRY).await;
assert_eq!(guest.set_up(), ["web"]);
shutdown.cancel();
task.await.unwrap();
}
}