use super::boot::boot_sandbox;
use super::cleanup::{inst_to_info, remove_sandbox_impl};
use super::persistence::{ProvisionIntent, SandboxProvisionOutcome, SandboxTransition};
use super::types::{SandboxBootTask, action};
use super::*;
impl SandboxManager {
pub async fn replay_sandbox_create(
&self,
id: &str,
request_key: &str,
) -> Result<Option<(SandboxId, String)>> {
self.await_reconcile().await?;
self.records
.replay_provision(id, request_key)
.map(|outcome| outcome.map(|outcome| (id.to_owned(), outcome.ip_address)))
}
pub async fn create_sandbox(&self, spec: SandboxSpec) -> Result<(SandboxId, String)> {
self.create_sandbox_keyed(spec, &Uuid::new_v4().to_string())
.await
}
pub async fn create_sandbox_keyed(
&self,
mut spec: SandboxSpec,
request_key: &str,
) -> Result<(SandboxId, String)> {
self.await_reconcile().await?;
let defaults = &self.config.defaults;
if spec.kernel.is_empty() {
spec.kernel.clone_from(&defaults.kernel);
}
if spec.rootfs.is_empty() {
spec.rootfs.clone_from(&defaults.rootfs);
}
if spec.boot_args.is_empty() {
spec.boot_args.clone_from(&defaults.boot_args);
}
if spec.vcpus == 0 {
spec.vcpus = defaults.vcpus as u32;
}
if spec.memory_mib == 0 {
spec.memory_mib = defaults.memory_mib;
}
if spec.network.mode.is_empty() {
spec.network.mode = "tap".into();
}
let id = spec
.id
.clone()
.filter(|s| !s.is_empty())
.unwrap_or_else(|| Uuid::new_v4().to_string());
super::validate_id("sandbox id", &id)?;
spec.id = Some(id.clone());
let vm_dir = PathBuf::from(&self.config.firecracker.data_dir)
.join("sandboxes")
.join(&id);
let reservation = super::reserve_id(
&self.instances,
&id,
SandboxInstance::new(id.clone(), spec.clone(), None, vm_dir.clone()),
)?;
let record = match self.records.provision_intent(&id, request_key, spec)? {
ProvisionIntent::Created(record) | ProvisionIntent::Resume(record) => record,
ProvisionIntent::Replay(record) => {
let outcome = record
.provision_outcome
.ok_or_else(|| VmmError::WrongState {
id: id.clone(),
expected: "a persisted create outcome".into(),
actual: "none".into(),
})?;
return Ok((id, outcome.ip_address));
}
ProvisionIntent::Blocked(_) => return Err(VmmError::AlreadyExists(id)),
};
let generation = record.generation;
let spec = record.effective_spec;
let arc = reservation.instance();
let mut creating_instance = arc.lock().unwrap();
creating_instance.record_generation = Some(generation);
creating_instance.labels.clone_from(&spec.labels);
creating_instance.spec.clone_from(&spec);
let mut net_alloc = None;
let setup = (|| -> Result<(String, Option<String>)> {
if spec.network.mode != "none" {
net_alloc = Some(self.network.reserve(&id)?);
}
let ip_address = net_alloc
.as_ref()
.map(|net| net.ip_address.to_string())
.unwrap_or_default();
super::reconcile::create_runtime_dir(&vm_dir)?;
let cleanup_record = super::reconcile::SandboxStateRecord::new(
&id,
None,
net_alloc.as_ref(),
None,
self.config.firecracker.jailer.is_some(),
None,
);
super::reconcile::write_state_record(&vm_dir, &cleanup_record)?;
if let Some(net) = &net_alloc {
self.network.activate(net)?;
}
let outcome = SandboxProvisionOutcome {
ip_address: ip_address.clone(),
};
let commit =
self.records
.transition(&id, generation, SandboxTransition::Starting(outcome))?;
Ok((ip_address, commit.durability_error))
})();
let (ip_address, starting_durability_error) = match setup {
Ok(result) => result,
Err(error) => {
let mut rollback_errors = Vec::new();
let mut network_cleanup_failed = false;
if let Some(net) = &net_alloc
&& let Err(release_error) = self.network.release_checked(net)
{
network_cleanup_failed = true;
rollback_errors.push(format!("network: {release_error}"));
}
if !network_cleanup_failed
&& let Err(remove_error) = std::fs::remove_dir_all(&vm_dir)
&& remove_error.kind() != std::io::ErrorKind::NotFound
{
rollback_errors.push(format!("directory {}: {remove_error}", vm_dir.display()));
}
if !rollback_errors.is_empty() {
creating_instance.network.clone_from(&net_alloc);
creating_instance.state = SandboxState::Failed;
creating_instance.error = Some(error.to_string());
let record_error = self
.records
.transition(
&id,
generation,
SandboxTransition::Failed(error.to_string()),
)
.err()
.map(|record_error| format!("record: {record_error}"));
if let Some(record_error) = record_error {
rollback_errors.push(record_error);
}
drop(creating_instance);
reservation.commit();
return Err(VmmError::Other(format!(
"{error}; sandbox rollback is incomplete: {}",
rollback_errors.join("; ")
)));
}
let abort = self.records.abort_provision(&id, generation)?;
if let Some(durability_error) = abort.durability_error {
return Err(VmmError::Unavailable(format!(
"{error}; create rollback is visible, but durability is unconfirmed: {durability_error}"
)));
}
return Err(error);
}
};
creating_instance.network.clone_from(&net_alloc);
let ttl_armed_for = Arc::downgrade(&arc);
{
let instances = Arc::clone(&self.instances);
let network = Arc::clone(&self.network);
let config = Arc::clone(&self.config);
let events_tx = self.events_tx.clone();
let cow_manager = Arc::clone(&self.cow_manager);
let records = Arc::clone(&self.records);
let id_clone = id.clone();
let spec_clone = spec.clone();
let net_alloc_clone = net_alloc;
let (resource_handoff_tx, resource_handoff) = tokio::sync::oneshot::channel();
let handle = tokio::spawn(async move {
boot_sandbox(
id_clone,
spec_clone,
net_alloc_clone,
vm_dir,
instances,
network,
config,
events_tx,
cow_manager,
records,
generation,
resource_handoff_tx,
)
.await;
});
creating_instance.boot_task = Some(SandboxBootTask {
resource_handoff: Some(resource_handoff),
handle,
});
}
drop(creating_instance);
reservation.commit();
let _ = self.events_tx.send(SandboxEvent::new(&id, action::CREATED));
if spec.ttl_seconds > 0 {
let instances = Arc::clone(&self.instances);
let network = Arc::clone(&self.network);
let events_tx = self.events_tx.clone();
let config2 = Arc::clone(&self.config);
let cow2 = Arc::clone(&self.cow_manager);
let records = Arc::clone(&self.records);
let id2 = id.clone();
let ttl = spec.ttl_seconds;
let armed_for = ttl_armed_for;
tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(ttl as u64)).await;
super::cleanup::expire_sandbox(
&id2,
Some(generation),
&armed_for,
&instances,
&network,
&events_tx,
&config2,
&cow2,
&records,
)
.await;
});
}
info!(sandbox_id = %id, "sandbox create requested (async boot started)");
if let Some(error) = starting_durability_error {
return Err(VmmError::Unavailable(format!(
"sandbox {id} was created, but ACK durability is unconfirmed: {error}"
)));
}
Ok((id, ip_address))
}
pub async fn stop_sandbox(&self, id: &SandboxId, timeout_seconds: u32) -> Result<()> {
self.await_reconcile().await?;
let budget = Duration::from_secs(u64::from(if timeout_seconds > 0 {
timeout_seconds
} else {
30
}));
let deadline = tokio::time::Instant::now() + budget;
let instance = self.get_instance(id)?;
let cleanup_lock = instance.lock().unwrap().cleanup_lock.clone();
let _cleanup_guard = cleanup_lock.lock().await;
super::ensure_current_instance(&self.instances, id, &instance)?;
let already_stopped = {
let inst = instance.lock().unwrap();
(inst.state == SandboxState::Stopped)
.then(|| (inst.record_generation, inst.vm_dir.clone()))
};
if let Some((generation, vm_dir)) = already_stopped {
if let Some(generation) = generation {
self.records
.transition(id, generation, SandboxTransition::Stopped)?
.confirmed("sandbox stop retry")?;
}
super::reconcile::clear_state_record(&vm_dir)?;
return Ok(());
}
let (was_running, vm_handle, record_generation, last_exited_at) = {
let mut inst = instance.lock().unwrap();
match inst.state {
SandboxState::Ready | SandboxState::Running | SandboxState::Stopping => {}
s => {
return Err(VmmError::WrongState {
id: id.clone(),
expected: "Ready, Running, or Stopping".into(),
actual: s.to_string(),
});
}
}
let was_running = inst.state == SandboxState::Running;
let captured = (
was_running,
inst.vm.as_ref().map(Arc::clone),
inst.record_generation,
inst.last_exited_at,
);
if let Some(generation) = inst.record_generation {
let commit =
self.records
.transition(id, generation, SandboxTransition::Stopping)?;
if let Some(error) = commit.durability_error {
warn!(
sandbox_id = %id,
error,
"stopping transition is visible but durability is unconfirmed"
);
}
}
inst.state = SandboxState::Stopping;
captured
};
let _ = self.events_tx.send(SandboxEvent::new(id, action::STOPPING));
if was_running {
while tokio::time::Instant::now() < deadline {
if instance.lock().unwrap().last_exited_at != last_exited_at {
break;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
if let Some(vm) = vm_handle {
let _ = tokio::time::timeout(Duration::from_secs(5), vm.send_ctrl_alt_del()).await;
}
let fc_process = instance.lock().unwrap().process.take();
if let Some(mut proc) = fc_process {
let remaining = deadline
.checked_duration_since(tokio::time::Instant::now())
.unwrap_or(Duration::from_secs(1))
.max(Duration::from_secs(1));
match tokio::time::timeout(remaining, proc.wait()).await {
Ok(Ok(_)) => {}
Ok(Err(error)) => {
instance.lock().unwrap().process = Some(proc);
return Err(VmmError::Process(format!(
"wait for sandbox {id} firecracker: {error}"
)));
}
Err(_) => {
warn!(sandbox_id = %id, "guest did not shut down in time; killing firecracker");
if let Err(error) = super::boot::kill_and_reap_fc_checked(&mut proc).await {
instance.lock().unwrap().process = Some(proc);
return Err(error);
}
}
}
}
let stop_commit = {
super::cleanup::release_runtime_resources(
id,
&instance,
&self.network,
&self.config,
&self.cow_manager,
)
.await?;
let commit = record_generation
.map(|generation| {
self.records
.transition(id, generation, SandboxTransition::Stopped)
})
.transpose()?;
let mut inst = instance.lock().unwrap();
inst.state = SandboxState::Stopped;
if commit
.as_ref()
.is_none_or(|commit| commit.durability_error.is_none())
{
super::reconcile::clear_state_record(&inst.vm_dir)?;
}
commit
};
let _ = self.events_tx.send(SandboxEvent::new(id, action::STOPPED));
info!(sandbox_id = %id, "sandbox stopped");
stop_commit
.map(|commit| commit.confirmed("sandbox stop"))
.transpose()?;
Ok(())
}
pub async fn remove_sandbox(&self, id: &SandboxId, force: bool) -> Result<()> {
self.await_reconcile().await?;
let expected = match self.get_instance(id) {
Ok(expected) => expected,
Err(VmmError::NotFound(_)) => {
let vm_dir = PathBuf::from(&self.config.firecracker.data_dir)
.join("sandboxes")
.join(id);
match super::reserve_id(
&self.instances,
id,
SandboxInstance::new(
id.clone(),
SandboxSpec {
id: Some(id.clone()),
..Default::default()
},
None,
vm_dir,
),
) {
Ok(_reservation) => {
let commit = self.records.cancel_pending_or_missing(id)?;
if let Some(error) = commit.durability_error {
return Err(VmmError::Unavailable(format!(
"sandbox {id} removal is visible, but durability is unconfirmed: {error}"
)));
}
info!(sandbox_id = %id, "sandbox already removed");
return Ok(());
}
Err(VmmError::AlreadyExists(_)) => self.get_instance(id)?,
Err(error) => return Err(error),
}
}
Err(error) => return Err(error),
};
remove_sandbox_impl(
id,
force,
&expected,
&self.instances,
&self.network,
&self.events_tx,
&self.config,
&self.cow_manager,
&self.records,
)
.await?;
info!(sandbox_id = %id, "sandbox removed");
Ok(())
}
pub fn inspect_sandbox(&self, id: &SandboxId) -> Result<SandboxInfo> {
let instance = self.get_instance(id)?;
let inst = instance.lock().unwrap();
Ok(inst_to_info(&inst))
}
pub fn list_sandboxes(
&self,
state_filter: Option<&str>,
label_filter: &HashMap<String, String>,
) -> Result<Vec<SandboxSummary>> {
self.check_reconcile()?;
let instances: Vec<_> = self.instances.read().unwrap().values().cloned().collect();
Ok(instances
.iter()
.filter_map(|arc| {
let inst = arc.lock().unwrap();
if let Some(sf) = state_filter
&& !sf.is_empty()
&& inst.state.to_string() != sf
{
return None;
}
for (k, v) in label_filter {
if inst.labels.get(k).map(String::as_str) != Some(v.as_str()) {
return None;
}
}
Some(SandboxSummary {
id: inst.id.clone(),
state: inst.state,
labels: inst.labels.clone(),
ip_address: inst
.network
.as_ref()
.map(|n| n.ip_address.to_string())
.unwrap_or_default(),
created_at: inst.created_at,
})
})
.collect())
}
pub fn subscribe_events(&self) -> broadcast::Receiver<SandboxEvent> {
self.events_tx.subscribe()
}
pub(super) fn get_instance(&self, id: &SandboxId) -> Result<Arc<Mutex<SandboxInstance>>> {
self.check_reconcile()?;
self.instances
.read()
.unwrap()
.get(id)
.cloned()
.ok_or_else(|| VmmError::NotFound(id.clone()))
}
pub(super) fn require_ready_vsock(&self, id: &SandboxId) -> Result<PathBuf> {
let instance = self.get_instance(id)?;
let inst = instance.lock().unwrap();
match inst.state {
SandboxState::Ready => {}
s => {
return Err(VmmError::WrongState {
id: id.clone(),
expected: "Ready".into(),
actual: s.to_string(),
});
}
}
inst.vsock_uds_path
.clone()
.ok_or_else(|| VmmError::Vsock(format!("sandbox {id} has no vsock configured")))
}
pub(super) fn require_alive_vsock(&self, id: &SandboxId) -> Result<PathBuf> {
let instance = self.get_instance(id)?;
let inst = instance.lock().unwrap();
match inst.state {
SandboxState::Ready | SandboxState::Running => {}
s => {
return Err(VmmError::WrongState {
id: id.clone(),
expected: "Ready or Running".into(),
actual: s.to_string(),
});
}
}
inst.vsock_uds_path
.clone()
.ok_or_else(|| VmmError::Vsock(format!("sandbox {id} has no vsock configured")))
}
pub async fn read_sandbox_file(&self, id: &SandboxId, path: &str) -> Result<Vec<u8>> {
let uds = self.require_alive_vsock(id)?;
crate::file_io::read_file(&uds, path).await
}
pub async fn write_sandbox_file(
&self,
id: &SandboxId,
path: &str,
mode: u32,
data: &[u8],
) -> Result<()> {
let uds = self.require_alive_vsock(id)?;
crate::file_io::write_file(&uds, path, mode, data).await
}
pub(super) fn get_vm_handle(&self, id: &SandboxId) -> Result<Arc<fc_sdk::Vm>> {
let instance = self.get_instance(id)?;
let inst = instance.lock().unwrap();
inst.vm
.as_ref()
.map(Arc::clone)
.ok_or_else(|| VmmError::WrongState {
id: id.clone(),
expected: "Ready or Running (VM handle not yet available)".into(),
actual: inst.state.to_string(),
})
}
}