mod tasks;
mod types;
pub(crate) use crate::litebox::box_impl::LiveState;
use crate::litebox::BoxStatus;
use crate::litebox::config::BoxConfig;
use crate::metrics::BoxMetricsStorage;
use crate::pipeline::{
BoxedTask, ExecutionPlan, PipelineBuilder, PipelineExecutor, PipelineMetrics, Stage,
};
use crate::runtime::rt_impl::SharedRuntimeImpl;
use crate::runtime::types::BoxState;
use boxlite_shared::errors::{BoxliteError, BoxliteResult};
use std::sync::Arc;
use tokio::sync::Mutex;
pub(crate) use tasks::PortPublishTask;
use tasks::{
BootAssetsTask, ContainerRootfsTask, FilesystemTask, GuestConnectTask, GuestInitTask,
GuestRootfsTask, InitCtx, VmmAttachTask, VmmSpawnTask,
};
use types::InitPipelineContext;
fn get_execution_plan(status: BoxStatus) -> BoxliteResult<ExecutionPlan<InitCtx>> {
let stages: Vec<Stage<BoxedTask<InitCtx>>> = match status {
BoxStatus::Configured => vec![
Stage::sequential(vec![Box::new(FilesystemTask)]),
Stage::parallel(vec![
Box::new(BootAssetsTask),
Box::new(ContainerRootfsTask),
Box::new(GuestRootfsTask),
]),
Stage::sequential(vec![Box::new(VmmSpawnTask)]),
Stage::sequential(vec![Box::new(GuestConnectTask)]),
Stage::parallel(vec![Box::new(GuestInitTask), Box::new(PortPublishTask)]),
],
BoxStatus::Stopped | BoxStatus::Failed => vec![
Stage::sequential(vec![Box::new(FilesystemTask)]),
Stage::parallel(vec![
Box::new(BootAssetsTask),
Box::new(ContainerRootfsTask),
Box::new(GuestRootfsTask),
]),
Stage::sequential(vec![Box::new(VmmSpawnTask)]),
Stage::sequential(vec![Box::new(GuestConnectTask)]),
Stage::parallel(vec![Box::new(GuestInitTask), Box::new(PortPublishTask)]),
],
BoxStatus::Running => vec![
Stage::sequential(vec![Box::new(VmmAttachTask)]),
Stage::sequential(vec![Box::new(GuestConnectTask)]),
Stage::sequential(vec![Box::new(PortPublishTask)]),
],
other => {
return Err(BoxliteError::InvalidState(format!(
"Cannot initialize box in {other} state"
)));
}
};
Ok(ExecutionPlan::new(stages))
}
fn box_metrics_from_pipeline(pipeline_metrics: &PipelineMetrics) -> BoxMetricsStorage {
let mut metrics = BoxMetricsStorage::new();
if let Some(duration_ms) = pipeline_metrics.task_duration_ms("filesystem_setup") {
metrics.set_stage_filesystem_setup(duration_ms);
}
if let Some(duration_ms) = pipeline_metrics.task_duration_ms("container_rootfs_prep") {
metrics.set_stage_image_prepare(duration_ms);
}
if let Some(duration_ms) = pipeline_metrics.task_duration_ms("guest_rootfs_init") {
metrics.set_stage_guest_rootfs(duration_ms);
}
if let Some(duration_ms) = pipeline_metrics.task_duration_ms("vmm_spawn") {
metrics.set_stage_box_spawn(duration_ms);
}
if let Some(duration_ms) = pipeline_metrics.task_duration_ms("vmm_attach") {
metrics.set_stage_box_spawn(duration_ms);
}
if let Some(_duration_ms) = pipeline_metrics.task_duration_ms("guest_connect") {
}
if let Some(duration_ms) = pipeline_metrics.task_duration_ms("guest_init") {
metrics.set_stage_container_init(duration_ms);
}
metrics
}
pub(crate) struct BoxBuilder {
runtime: SharedRuntimeImpl,
config: BoxConfig,
state: BoxState,
}
fn validate_persisted_options(
features: &crate::experimental::ExperimentalFeatures,
options: &crate::BoxOptions,
) -> BoxliteResult<()> {
features.require_for_options(options)?;
options.sanitize_persisted()
}
impl BoxBuilder {
pub(crate) fn new(
runtime: SharedRuntimeImpl,
config: BoxConfig,
state: BoxState,
) -> BoxliteResult<Self> {
let options = &config.options;
validate_persisted_options(&runtime.experimental_features, options)?;
Ok(Self {
runtime,
config,
state,
})
}
pub(crate) async fn build(self) -> BoxliteResult<(LiveState, types::CleanupGuard)> {
use std::time::Instant;
let total_start = Instant::now();
let BoxBuilder {
runtime,
config,
state,
} = self;
let status = state.status;
let reuse_rootfs = status == BoxStatus::Stopped;
let skip_guest_wait = status == BoxStatus::Running;
let ctx = InitPipelineContext::new(config, runtime.clone(), reuse_rootfs, skip_guest_wait);
let ctx = Arc::new(Mutex::new(ctx));
let ctx_for_cleanup = Arc::clone(&ctx);
let inner = async move {
let plan = get_execution_plan(status)?;
let pipeline = PipelineBuilder::from_plan(plan);
let pipeline_metrics = PipelineExecutor::execute(pipeline, Arc::clone(&ctx)).await?;
let mut ctx = ctx.lock().await;
let total_create_duration_ms = total_start.elapsed().as_millis();
let handler = ctx
.guard
.take_handler()
.ok_or_else(|| BoxliteError::Internal("handler was not set".into()))?;
let mut metrics = box_metrics_from_pipeline(&pipeline_metrics);
metrics.set_total_create_duration(total_create_duration_ms);
metrics.log_init_stages();
let guest_session = ctx.guest_session.take().ok_or_else(|| {
BoxliteError::Internal("guest_connect task must run first".into())
})?;
let (container_disk, guest_disk) = if status == BoxStatus::Running {
use crate::disk::DiskFormat;
use crate::disk::constants::filenames;
let disk = crate::disk::Disk::new(
ctx.config.box_home.join(filenames::CONTAINER_DISK),
DiskFormat::Qcow2,
true,
);
(disk, None)
} else {
let container_disk = ctx
.container_disk
.take()
.ok_or_else(|| BoxliteError::Internal("rootfs task must run first".into()))?;
(container_disk, ctx.guest_disk.take())
};
#[cfg(target_os = "linux")]
let bind_mount = ctx.bind_mount.take();
let mut placeholder =
types::CleanupGuard::new(ctx.runtime.clone(), ctx.config.id.clone());
placeholder.disarm();
let guard = std::mem::replace(&mut ctx.guard, placeholder);
let network = ctx.network_backend.take();
let published_ports = ctx.published_ports.take();
let live_state = LiveState::new(
handler,
guest_session,
network,
published_ports,
metrics,
container_disk,
guest_disk,
#[cfg(target_os = "linux")]
bind_mount,
);
Ok::<(LiveState, types::CleanupGuard), BoxliteError>((live_state, guard))
};
match inner.await {
Ok(parts) => Ok(parts),
Err(e) => {
ctx_for_cleanup.lock().await.guard.set_last_error(&e);
Err(e)
}
}
}
}
#[cfg(test)]
mod plan_tests {
use super::*;
use crate::experimental::ExperimentalFeatures;
use crate::experimental::custom_kernel::{KernelFormat, KernelOptions};
use crate::pipeline::ExecutionMode;
fn plan_shape(status: BoxStatus) -> Vec<(ExecutionMode, Vec<String>)> {
get_execution_plan(status)
.unwrap()
.stages()
.into_iter()
.map(|stage| {
(
stage.execution,
stage
.tasks
.into_iter()
.map(|task| task.name().to_string())
.collect(),
)
})
.collect()
}
fn boot_pipeline_shape() -> Vec<(ExecutionMode, Vec<String>)> {
vec![
(
ExecutionMode::Sequential,
vec!["filesystem_setup".to_string()],
),
(
ExecutionMode::Parallel,
vec![
"boot_assets_prepare".to_string(),
"container_rootfs_prep".to_string(),
"guest_rootfs_init".to_string(),
],
),
(ExecutionMode::Sequential, vec!["vmm_spawn".to_string()]),
(ExecutionMode::Sequential, vec!["guest_connect".to_string()]),
(
ExecutionMode::Parallel,
vec!["guest_init".to_string(), "port_publish".to_string()],
),
]
}
#[test]
fn configured_box_prepares_boot_assets_before_vmm_spawn() {
assert_eq!(plan_shape(BoxStatus::Configured), boot_pipeline_shape());
}
#[test]
fn restartable_boxes_prepare_boot_assets_before_vmm_spawn() {
let expected = boot_pipeline_shape();
assert_eq!(plan_shape(BoxStatus::Stopped), expected);
assert_eq!(plan_shape(BoxStatus::Failed), boot_pipeline_shape());
}
#[test]
fn running_box_reattach_does_not_prepare_boot_assets() {
assert_eq!(
plan_shape(BoxStatus::Running),
vec![
(ExecutionMode::Sequential, vec!["vmm_attach".to_string()]),
(ExecutionMode::Sequential, vec!["guest_connect".to_string()]),
(ExecutionMode::Sequential, vec!["port_publish".to_string()]),
]
);
}
#[test]
#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))]
fn persisted_custom_kernel_uses_injected_feature_state() {
#[cfg(target_arch = "x86_64")]
let format = KernelFormat::Elf;
#[cfg(target_arch = "aarch64")]
let format = KernelFormat::PeGz;
let mut options = crate::BoxOptions::default();
options.advanced.kernel =
Some(KernelOptions::new("/source/no-longer-needed").with_format(format));
let error = validate_persisted_options(&ExperimentalFeatures::default(), &options)
.expect_err("persisted custom kernel must be disabled by default");
assert!(
error
.to_string()
.contains("ExperimentalFeature::CustomKernel")
);
let enabled = ExperimentalFeatures::parse("custom-kernel").unwrap();
validate_persisted_options(&enabled, &options).unwrap();
}
#[test]
fn persisted_nested_virtualization_uses_injected_feature_state() {
let mut advanced = crate::runtime::advanced_options::AdvancedBoxOptions::default();
advanced.nested_virtualization = true;
let options = crate::BoxOptions {
advanced,
..Default::default()
};
let error = validate_persisted_options(&ExperimentalFeatures::default(), &options)
.expect_err("persisted nested virtualization must be disabled by default");
assert!(
error
.to_string()
.contains("ExperimentalFeature::NestedVirtualization")
);
let enabled = ExperimentalFeatures::parse("nested-virtualization").unwrap();
validate_persisted_options(&enabled, &options).unwrap();
}
}