#![cfg(unix)]
use std::collections::HashMap;
use std::path::Path;
use containerd_client;
use containerd_client::services::v1::containers_client::ContainersClient;
use containerd_client::services::v1::content_client::ContentClient;
use containerd_client::services::v1::images_client::ImagesClient;
use containerd_client::services::v1::leases_client::LeasesClient;
use containerd_client::services::v1::{
Container, DeleteContentRequest, GetContainerRequest, GetImageRequest, Image, Info,
InfoRequest, ReadContentRequest, UpdateRequest, WriteAction, WriteContentRequest,
WriteContentResponse,
};
use containerd_client::tonic::transport::Channel;
use containerd_client::tonic::Streaming;
use containerd_client::{tonic, with_namespace};
use futures::TryStreamExt;
use oci_spec::image::{Arch, ImageManifest, MediaType, Platform};
use prost_types::FieldMask;
use sha256::digest;
use tokio::runtime::Runtime;
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use tonic::{Code, Request};
use super::lease::LeaseGuard;
use crate::container::Engine;
use crate::sandbox::error::{Error as ShimError, Result};
use crate::sandbox::oci::{self, WasmLayer};
use crate::with_lease;
static PRECOMPILE_PREFIX: &str = "runwasi.io/precompiled";
static MAX_WRITE_CHUNK_SIZE_BYTES: i64 = 1024 * 1024 * 15;
pub struct Client {
inner: Channel,
rt: Runtime,
namespace: String,
address: String,
}
#[derive(Debug)]
pub(crate) struct WriteContent {
_lease: LeaseGuard,
pub digest: String,
}
impl Client {
pub fn connect(
address: impl AsRef<Path> + ToString,
namespace: impl ToString,
) -> Result<Client> {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let inner = rt
.block_on(containerd_client::connect(address.as_ref()))
.map_err(|err| ShimError::Containerd(err.to_string()))?;
Ok(Client {
inner,
rt,
namespace: namespace.to_string(),
address: address.to_string(),
})
}
fn read_content(&self, digest: impl ToString) -> Result<Vec<u8>> {
self.rt.block_on(async {
let req = ReadContentRequest {
digest: digest.to_string(),
..Default::default()
};
let req = with_namespace!(req, self.namespace);
ContentClient::new(self.inner.clone())
.read(req)
.await
.map_err(|err| ShimError::Containerd(err.to_string()))?
.into_inner()
.map_ok(|msg| msg.data)
.try_concat()
.await
.map_err(|err| ShimError::Containerd(err.to_string()))
})
}
#[allow(dead_code)]
fn delete_content(&self, digest: impl ToString) -> Result<()> {
self.rt.block_on(async {
let req = DeleteContentRequest {
digest: digest.to_string(),
};
let req = with_namespace!(req, self.namespace);
ContentClient::new(self.inner.clone())
.delete(req)
.await
.map_err(|err| ShimError::Containerd(err.to_string()))?;
Ok(())
})
}
fn lease(&self, reference: String) -> Result<LeaseGuard> {
self.rt.block_on(async {
let mut lease_labels = HashMap::new();
let expire = chrono::Utc::now() + chrono::Duration::try_hours(24).unwrap();
lease_labels.insert("containerd.io/gc.expire".to_string(), expire.to_rfc3339());
let lease_request = containerd_client::services::v1::CreateRequest {
id: reference.clone(),
labels: lease_labels,
};
let mut leases_client = LeasesClient::new(self.inner.clone());
let lease = leases_client
.create(with_namespace!(lease_request, self.namespace))
.await
.map_err(|e| ShimError::Containerd(e.to_string()))?
.into_inner()
.lease
.ok_or_else(|| {
ShimError::Containerd(format!("unable to create lease for {}", reference))
})?;
Ok(LeaseGuard {
lease_id: lease.id,
address: self.address.clone(),
namespace: self.namespace.clone(),
})
})
}
fn save_content(
&self,
data: Vec<u8>,
unique_id: &str,
labels: HashMap<String, String>,
) -> Result<WriteContent> {
let expected = format!("sha256:{}", digest(data.clone()));
let reference = format!("precompile-{}", unique_id);
let lease = self.lease(reference.clone())?;
let digest = self.rt.block_on(async {
let (tx, rx) = mpsc::channel(1);
let len = data.len() as i64;
log::debug!("Writing {} bytes to content store", len);
let mut client = ContentClient::new(self.inner.clone());
let req = WriteContentRequest {
r#ref: reference.clone(),
action: WriteAction::Stat.into(),
expected: expected.clone(),
..Default::default()
};
tx.send(req)
.await
.map_err(|err| ShimError::Containerd(err.to_string()))?;
let request_stream = ReceiverStream::new(rx);
let request_stream =
with_lease!(request_stream, self.namespace, lease.lease_id.clone());
let mut response_stream = match client.write(request_stream).await {
Ok(response_stream) => response_stream.into_inner(),
Err(e) if e.code() == Code::AlreadyExists => {
log::info!("content already exists {}", expected.clone().to_string());
return Ok(expected);
}
Err(e) => return Err(ShimError::Containerd(e.to_string())),
};
let response = response_stream
.message()
.await
.map_err(|e| ShimError::Containerd(e.to_string()))?
.ok_or_else(|| {
ShimError::Containerd(format!(
"no response received after write request for {}",
expected
))
})?;
log::debug!(
"Starting to write content for layer {} with current status response {:?}",
expected,
response
);
let mut offset = response.offset;
while offset < len {
let end = (offset + MAX_WRITE_CHUNK_SIZE_BYTES).min(len);
let chunk = &data[offset as usize..end as usize];
let write_request = WriteContentRequest {
action: WriteAction::Write.into(),
total: 0,
offset,
data: chunk.to_vec(),
..Default::default()
};
let response =
send_message(write_request, &mut response_stream, &tx, &expected).await?;
log::debug!(
"Writing content for layer {} at offset {} got response: {:?}",
expected,
offset,
response
);
offset = end;
}
let commit_request = WriteContentRequest {
action: WriteAction::Commit.into(),
total: len,
offset: len,
expected: expected.clone(),
labels,
data: Vec::new(),
..Default::default()
};
let response =
send_message(commit_request, &mut response_stream, &tx, &expected).await?;
log::info!(
"Validating final response after writing content for layer {}: {:?}",
expected,
response
);
if response.offset != len {
return Err(ShimError::Containerd(format!(
"failed to write all bytes, expected {} got {}",
len, response.offset
)));
}
if response.digest != expected {
return Err(ShimError::Containerd(format!(
"unexpected digest, expected {} got {}",
expected, response.digest
)));
}
Ok(response.digest)
})?;
Ok(WriteContent {
_lease: lease,
digest: digest.clone(),
})
}
fn get_info(&self, content_digest: &str) -> Result<Info> {
self.rt.block_on(async {
let req = InfoRequest {
digest: content_digest.to_string(),
};
let req = with_namespace!(req, self.namespace);
let info = ContentClient::new(self.inner.clone())
.info(req)
.await
.map_err(|err| ShimError::Containerd(err.to_string()))?
.into_inner()
.info
.ok_or_else(|| {
ShimError::Containerd(format!(
"failed to get info for content {}",
content_digest
))
})?;
Ok(info)
})
}
fn update_info(&self, info: Info) -> Result<Info> {
self.rt.block_on(async {
let req = UpdateRequest {
info: Some(info.clone()),
update_mask: Some(FieldMask {
paths: vec!["labels".to_string()],
}),
};
let req = with_namespace!(req, self.namespace);
let info = ContentClient::new(self.inner.clone())
.update(req)
.await
.map_err(|err| ShimError::Containerd(err.to_string()))?
.into_inner()
.info
.ok_or_else(|| {
ShimError::Containerd(format!(
"failed to update info for content {}",
info.digest
))
})?;
Ok(info)
})
}
fn get_image(&self, image_name: impl ToString) -> Result<Image> {
self.rt.block_on(async {
let name = image_name.to_string();
let req = GetImageRequest { name };
let req = with_namespace!(req, self.namespace);
let image = ImagesClient::new(self.inner.clone())
.get(req)
.await
.map_err(|err| ShimError::Containerd(err.to_string()))?
.into_inner()
.image
.ok_or_else(|| {
ShimError::Containerd(format!(
"failed to get image for image {}",
image_name.to_string()
))
})?;
Ok(image)
})
}
fn extract_image_content_sha(&self, image: &Image) -> Result<String> {
let digest = image
.target
.as_ref()
.ok_or_else(|| {
ShimError::Containerd(format!(
"failed to get image content sha for image {}",
image.name
))
})?
.digest
.clone();
Ok(digest)
}
fn get_container(&self, container_name: impl ToString) -> Result<Container> {
self.rt.block_on(async {
let id = container_name.to_string();
let req = GetContainerRequest { id };
let req = with_namespace!(req, self.namespace);
let container = ContainersClient::new(self.inner.clone())
.get(req)
.await
.map_err(|err| ShimError::Containerd(err.to_string()))?
.into_inner()
.container
.ok_or_else(|| {
ShimError::Containerd(format!(
"failed to get image for container {}",
container_name.to_string()
))
})?;
Ok(container)
})
}
fn get_image_manifest_and_digest(&self, image_name: &str) -> Result<(ImageManifest, String)> {
let image = self.get_image(image_name)?;
let image_digest = self.extract_image_content_sha(&image)?;
let manifest = ImageManifest::from_reader(self.read_content(&image_digest)?.as_slice())?;
Ok((manifest, image_digest))
}
pub fn load_modules<T: Engine>(
&self,
containerd_id: impl ToString,
engine: &T,
) -> Result<(Vec<oci::WasmLayer>, Platform)> {
let container = self.get_container(containerd_id.to_string())?;
let (manifest, image_digest) = self.get_image_manifest_and_digest(&container.image)?;
let image_config_descriptor = manifest.config();
let image_config = self.read_content(image_config_descriptor.digest())?;
let image_config = image_config.as_slice();
let platform: Platform = serde_json::from_slice(image_config)?;
let Arch::Wasm = platform.architecture() else {
log::info!("manifest is not in WASM OCI image format");
return Ok((vec![], platform));
};
log::info!("found manifest with WASM OCI image format");
let (can_precompile, precompile_id) = match engine.can_precompile() {
Some(precompile_id) => (true, precompile_label(T::name(), &precompile_id)),
None => (false, "".to_string()),
};
let image_info = self.get_info(&image_digest)?;
let mut needs_precompile =
can_precompile && !image_info.labels.contains_key(&precompile_id);
let layers = manifest
.layers()
.iter()
.filter(|x| is_wasm_layer(x.media_type(), T::supported_layers_types()))
.map(|original_config| {
self.read_wasm_layer(
original_config,
can_precompile,
&precompile_id,
&mut needs_precompile,
)
})
.collect::<Result<Vec<_>>>()?;
if layers.is_empty() {
log::info!("no WASM layers found in OCI image");
return Ok((vec![], platform));
}
if needs_precompile {
log::info!("precompiling layers for image: {}", container.image);
let compiled_layers = match engine.precompile(&layers) {
Ok(compiled_layers) => {
if compiled_layers.len() != layers.len() {
return Err(ShimError::FailedPrecondition(
"precompile returned wrong number of layers".to_string(),
));
}
compiled_layers
}
Err(e) => {
log::error!("precompilation failed: {}", e);
return Ok((layers, platform));
}
};
let mut layers_for_runtime = Vec::with_capacity(compiled_layers.len());
for (i, compiled_layer) in compiled_layers.iter().enumerate() {
if compiled_layer.is_none() {
log::debug!("no compiled layer using original");
layers_for_runtime.push(layers[i].clone());
continue;
}
let compiled_layer = compiled_layer.as_ref().unwrap();
let original_config = &layers[i].config;
let labels = HashMap::from([(
format!("{precompile_id}/original"),
original_config.digest().to_string(),
)]);
let precompiled_content =
self.save_content(compiled_layer.clone(), &precompile_id, labels)?;
log::debug!(
"updating original layer {} with compiled layer {}",
original_config.digest(),
precompiled_content.digest
);
let mut original_layer = self.get_info(original_config.digest())?;
original_layer
.labels
.insert(precompile_id.clone(), precompiled_content.digest.clone());
original_layer.labels.insert(
format!("containerd.io/gc.ref.content.precompile.{}", i),
precompiled_content.digest.clone(),
);
self.update_info(original_layer)?;
log::debug!(
"updating image content with precompile digest to avoid garbage collection"
);
let mut image_content = self.get_info(&image_digest)?;
image_content.labels.insert(
format!("containerd.io/gc.ref.content.precompile.{}", i),
precompiled_content.digest,
);
image_content
.labels
.insert(precompile_id.clone(), "true".to_string());
self.update_info(image_content)?;
layers_for_runtime.push(WasmLayer {
config: original_config.clone(),
layer: compiled_layer.clone(),
});
}
return Ok((layers_for_runtime, platform));
};
log::info!("using OCI layers");
Ok((layers, platform))
}
fn read_wasm_layer(
&self,
original_config: &oci_spec::image::Descriptor,
can_precompile: bool,
precompile_id: &String,
needs_precompile: &mut bool,
) -> std::prelude::v1::Result<WasmLayer, ShimError> {
let mut digest_to_load = original_config.digest().clone();
if can_precompile {
let info = self.get_info(&digest_to_load)?;
if let Some(label) = info.labels.get(precompile_id) {
digest_to_load = label.clone();
log::info!(
"layer {} has pre-compiled content: {} ",
info.digest,
&digest_to_load
);
}
}
log::debug!("loading digest: {} ", &digest_to_load);
self.read_content(&digest_to_load)
.map(|module| WasmLayer {
config: original_config.clone(),
layer: module,
})
.or_else(|e| {
if digest_to_load != *original_config.digest() {
log::error!("failed to load precompiled layer: {}", e);
log::error!("falling back to original layer and marking for recompile");
*needs_precompile = can_precompile; self.read_content(original_config.digest())
.map(|module| WasmLayer {
config: original_config.clone(),
layer: module,
})
} else {
Err(e)
}
})
}
}
fn precompile_label(name: &str, version: &str) -> String {
format!("{}/{}/{}", PRECOMPILE_PREFIX, name, version)
}
fn is_wasm_layer(media_type: &MediaType, supported_layer_types: &[&str]) -> bool {
let supported = supported_layer_types.contains(&media_type.to_string().as_str());
log::debug!(
"layer type {} is supported: {}",
media_type.to_string().as_str(),
supported
);
supported
}
async fn send_message(
request: WriteContentRequest,
response_stream: &mut Streaming<WriteContentResponse>,
tx: &mpsc::Sender<WriteContentRequest>,
digest: &str,
) -> Result<WriteContentResponse> {
tx.send(request)
.await
.map_err(|err| ShimError::Containerd(format!("commit request error: {}", err)))?;
response_stream
.message()
.await
.map_err(|err| ShimError::Containerd(format!("response stream error: {}", err)))?
.ok_or_else(|| {
ShimError::Containerd(format!(
"no response received after write content request for {}",
digest
))
})
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use std::sync::atomic::{AtomicI32, Ordering};
use std::sync::Arc;
use oci_tar_builder::WASM_LAYER_MEDIA_TYPE;
use rand::prelude::*;
use super::*;
use crate::container::RuntimeContext;
use crate::sandbox::Stdio;
use crate::testing::oci_helpers::ImageContent;
use crate::testing::{oci_helpers, TEST_NAMESPACE};
#[test]
fn test_save_content() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, "test-ns").unwrap();
let data = b"hello world".to_vec();
let expected = digest(data.clone());
let expected = format!("sha256:{}", expected);
let label = HashMap::from([(precompile_label("test", "hasdfh"), "original".to_string())]);
let returned = client.save_content(data, "test", label.clone()).unwrap();
assert_eq!(expected, returned.digest.clone());
let data = client.read_content(returned.digest.clone()).unwrap();
assert_eq!(data, b"hello world");
client
.save_content(data.clone(), "test", label.clone())
.expect_err("Should not be able to save when lease is open");
drop(returned);
let returned = client.save_content(data, "test", label.clone()).unwrap();
assert_eq!(expected, returned.digest);
client.delete_content(expected.clone()).unwrap();
client
.read_content(expected)
.expect_err("content should not exist");
}
#[test]
fn test_layers_when_precompile_not_supported() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let (_, container_name, _cleanup) = generate_test_container(None, &[&fake_bytes]);
let engine = FakePrecomiplerEngine::new(None);
let (layers, _) = client.load_modules(container_name, &engine).unwrap();
assert_eq!(layers.len(), 1);
assert_eq!(layers[0].layer, fake_bytes.bytes);
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 0);
}
#[test]
fn test_layers_are_precompiled_once() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, crate::testing::TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let (_image_name, container_name, _cleanup) = generate_test_container(None, &[&fake_bytes]);
let fake_precompiled_bytes = generate_content("precompiled", WASM_LAYER_MEDIA_TYPE);
let mut engine = FakePrecomiplerEngine::new(Some(()));
engine.add_precompiled_bits(fake_bytes.bytes.clone(), &fake_precompiled_bytes);
let (_, _) = client.load_modules(&container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
let (layers, _) = client.load_modules(&container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
assert_eq!(layers.len(), 1);
assert_eq!(layers[0].layer, fake_precompiled_bytes.bytes);
}
#[test]
fn test_layers_are_recompiled_if_version_changes() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, crate::testing::TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let (_image_name, container_name, _cleanup) = generate_test_container(None, &[&fake_bytes]);
let fake_precompiled_bytes = generate_content("precompiled", WASM_LAYER_MEDIA_TYPE);
let mut engine = FakePrecomiplerEngine::new(Some(()));
engine.add_precompiled_bits(fake_bytes.bytes.clone(), &fake_precompiled_bytes);
let (_, _) = client.load_modules(&container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
engine.precompile_id = Some("new_version".to_string());
let (_, _) = client.load_modules(&container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 2);
}
#[test]
fn test_layers_are_precompiled() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, crate::testing::TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let (image_name, container_name, _cleanup) = generate_test_container(None, &[&fake_bytes]);
let fake_precompiled_bytes = generate_content("precompiled", WASM_LAYER_MEDIA_TYPE);
let mut engine = FakePrecomiplerEngine::new(Some(()));
engine.add_precompiled_bits(fake_bytes.bytes.clone(), &fake_precompiled_bytes);
let expected_id = precompile_label(
FakePrecomiplerEngine::name(),
engine.can_precompile().unwrap().as_str(),
);
let (layers, _) = client.load_modules(container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
assert_eq!(layers.len(), 1);
assert_eq!(layers[0].layer, fake_precompiled_bytes.bytes);
let (manifest, _) = client.get_image_manifest_and_digest(&image_name).unwrap();
let original_config = manifest.layers().first().unwrap();
let info = client.get_info(original_config.digest()).unwrap();
let actual_digest = info.labels.get(&expected_id).unwrap();
assert_eq!(
actual_digest.to_string(),
format!("sha256:{}", &digest(fake_precompiled_bytes.bytes.clone()))
);
}
#[test]
fn test_layers_are_precompiled_but_not_for_all_layers() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, crate::testing::TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let non_wasm_bytes = generate_content("original_dont_compile", "textfile");
let (_image_name, container_name, _cleanup) =
generate_test_container(None, &[&fake_bytes, &non_wasm_bytes]);
let fake_precompiled_bytes = generate_content("precompiled", WASM_LAYER_MEDIA_TYPE);
let mut engine = FakePrecomiplerEngine::new(Some(()));
engine.add_precompiled_bits(fake_bytes.bytes.clone(), &fake_precompiled_bytes);
let (layers, _) = client.load_modules(container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
assert_eq!(engine.layers_compiled_per_call.load(Ordering::SeqCst), 1);
assert_eq!(layers.len(), 2);
assert_eq!(layers[0].layer, fake_precompiled_bytes.bytes);
assert_eq!(layers[1].layer, non_wasm_bytes.bytes);
}
#[test]
fn test_layers_do_not_need_precompiled_if_new_layers_are_added_to_existing_image() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, crate::testing::TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let (_image_name, container_name, _cleanup) = generate_test_container(None, &[&fake_bytes]);
let fake_precompiled_bytes = generate_content("precompiled", WASM_LAYER_MEDIA_TYPE);
let mut engine = FakePrecomiplerEngine::new(Some(()));
engine.add_precompiled_bits(fake_bytes.bytes.clone(), &fake_precompiled_bytes);
let (layers, _) = client.load_modules(container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
assert_eq!(layers.len(), 1);
assert_eq!(layers[0].layer, fake_precompiled_bytes.bytes);
let image_sha = client
.get_image(&_image_name)
.unwrap()
.target
.unwrap()
.digest;
let fake_bytes2 = generate_content("image2", WASM_LAYER_MEDIA_TYPE);
let (_image_name2, container_name2, _cleanup2) =
generate_test_container(Some(_image_name), &[&fake_bytes, &fake_bytes2]);
let fake_precompiled_bytes2 = generate_content("precompiled2", WASM_LAYER_MEDIA_TYPE);
engine.add_precompiled_bits(fake_bytes2.bytes.clone(), &fake_precompiled_bytes2);
oci_helpers::wait_for_content_removal(&image_sha).unwrap();
let (layers, _) = client.load_modules(container_name2, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 2);
assert_eq!(layers.len(), 2);
assert_eq!(engine.layers_compiled_per_call.load(Ordering::SeqCst), 1);
}
#[test]
fn test_layers_do_not_need_precompiled_if_new_layers_are_add_to_new_image() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, crate::testing::TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let (_image_name, container_name, _cleanup) = generate_test_container(None, &[&fake_bytes]);
let fake_precompiled_bytes = generate_content("precompiled", WASM_LAYER_MEDIA_TYPE);
let mut engine = FakePrecomiplerEngine::new(Some(()));
engine.add_precompiled_bits(fake_bytes.bytes.clone(), &fake_precompiled_bytes);
let (layers, _) = client.load_modules(container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
assert_eq!(layers.len(), 1);
assert_eq!(layers[0].layer, fake_precompiled_bytes.bytes);
let fake_bytes2 = generate_content("image2", WASM_LAYER_MEDIA_TYPE);
let (_image_name2, container_name2, _cleanup2) =
generate_test_container(None, &[&fake_bytes, &fake_bytes2]);
let fake_precompiled_bytes2 = generate_content("precompiled2", WASM_LAYER_MEDIA_TYPE);
engine.add_precompiled_bits(fake_bytes2.bytes.clone(), &fake_precompiled_bytes2);
let (layers, _) = client.load_modules(container_name2, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 2);
assert_eq!(layers.len(), 2);
assert_eq!(engine.layers_compiled_per_call.load(Ordering::SeqCst), 1);
}
#[test]
fn test_layers_are_precompiled_for_multiple_layers() {
let path = PathBuf::from("/run/containerd/containerd.sock");
let path = path.to_str().unwrap();
let client = Client::connect(path, crate::testing::TEST_NAMESPACE).unwrap();
let fake_bytes = generate_content("original", WASM_LAYER_MEDIA_TYPE);
let fake_bytes2 = generate_content("original1", WASM_LAYER_MEDIA_TYPE);
let (image_name, container_name, _cleanup) =
generate_test_container(None, &[&fake_bytes, &fake_bytes2]);
let fake_precompiled_bytes = generate_content("precompiled", WASM_LAYER_MEDIA_TYPE);
let fake_precompiled_bytes2 = generate_content("precompiled1", WASM_LAYER_MEDIA_TYPE);
let mut engine = FakePrecomiplerEngine::new(Some(()));
engine.add_precompiled_bits(fake_bytes.bytes.clone(), &fake_precompiled_bytes);
engine.add_precompiled_bits(fake_bytes2.bytes.clone(), &fake_precompiled_bytes2);
let expected_id = precompile_label(
FakePrecomiplerEngine::name(),
engine.can_precompile().unwrap().as_str(),
);
let (layers, _) = client.load_modules(container_name, &engine).unwrap();
assert_eq!(engine.precompile_called.load(Ordering::SeqCst), 1);
assert_eq!(engine.layers_compiled_per_call.load(Ordering::SeqCst), 2);
assert_eq!(layers.len(), 2);
assert_eq!(layers[0].layer, fake_precompiled_bytes.bytes);
assert_eq!(layers[1].layer, fake_precompiled_bytes2.bytes);
let (manifest, _) = client.get_image_manifest_and_digest(&image_name).unwrap();
let original_config1 = manifest.layers().first().unwrap();
let info1 = client.get_info(original_config1.digest()).unwrap();
let actual_digest1 = info1.labels.get(&expected_id).unwrap();
assert_eq!(
actual_digest1.to_string(),
format!("sha256:{}", &digest(fake_precompiled_bytes.bytes.clone()))
);
let original_config2 = manifest.layers().last().unwrap();
let info2 = client.get_info(original_config2.digest()).unwrap();
let actual_digest2 = info2.labels.get(&expected_id).unwrap();
assert_eq!(
actual_digest2.to_string(),
format!("sha256:{}", &digest(fake_precompiled_bytes2.bytes.clone()))
);
}
fn generate_test_container(
name: Option<String>,
original: &[&oci_helpers::ImageContent],
) -> (String, String, oci_helpers::OCICleanup) {
let _ = env_logger::try_init();
let random_number = random_number();
let image_name = name.unwrap_or(format!("localhost/test:latest{}", random_number));
oci_helpers::import_image(&image_name, original).unwrap();
let container_name = format!("test-container-{}", random_number);
oci_helpers::create_container(&container_name, &image_name).unwrap();
let _cleanup = oci_helpers::OCICleanup {
image_name: image_name.clone(),
container_name: container_name.clone(),
};
(image_name, container_name, _cleanup)
}
fn generate_content(seed: &str, media_type: &str) -> oci_helpers::ImageContent {
let mut content = seed.as_bytes().to_vec();
for _ in 0..100 {
content.push(random_number() as u8);
}
ImageContent {
bytes: content,
media_type: media_type.to_string(),
}
}
fn random_number() -> u32 {
let x: u32 = random();
x
}
#[derive(Clone)]
struct FakePrecomiplerEngine {
precompile_id: Option<String>,
precompiled_layers: HashMap<String, Vec<u8>>,
precompile_called: Arc<AtomicI32>,
layers_compiled_per_call: Arc<AtomicI32>,
}
impl FakePrecomiplerEngine {
fn new(can_precompile: Option<()>) -> Self {
let precompile_id = match can_precompile {
Some(_) => {
let precompile_id = format!("uuid-{}", random_number());
Some(precompile_id)
}
None => None,
};
FakePrecomiplerEngine {
precompile_id,
precompiled_layers: HashMap::new(),
precompile_called: Arc::new(AtomicI32::new(0)),
layers_compiled_per_call: Arc::new(AtomicI32::new(0)),
}
}
fn add_precompiled_bits(
&mut self,
original: Vec<u8>,
precompiled_content: &oci_helpers::ImageContent,
) {
let key = digest(original);
self.precompiled_layers
.insert(key, precompiled_content.bytes.clone());
}
}
impl Engine for FakePrecomiplerEngine {
fn name() -> &'static str {
"fake"
}
fn run_wasi(
&self,
_ctx: &impl RuntimeContext,
_stdio: Stdio,
) -> std::result::Result<i32, anyhow::Error> {
panic!("not implemented")
}
fn can_precompile(&self) -> Option<String> {
self.precompile_id.clone()
}
fn supported_layers_types() -> &'static [&'static str] {
&[WASM_LAYER_MEDIA_TYPE, "textfile"]
}
fn precompile(&self, layers: &[WasmLayer]) -> Result<Vec<Option<Vec<u8>>>, anyhow::Error> {
self.layers_compiled_per_call.store(0, Ordering::SeqCst);
self.precompile_called.fetch_add(1, Ordering::SeqCst);
let mut compiled_layers = vec![];
for layer in layers {
if layer.config.media_type().to_string() == *"textfile" {
compiled_layers.push(None);
continue;
}
let key = digest(layer.layer.clone());
if self.precompiled_layers.values().any(|l| digest(l) == key) {
compiled_layers.push(None);
continue;
}
self.precompiled_layers.iter().all(|x| {
log::warn!("layer: {:?}", x.0);
true
});
let precompiled = self.precompiled_layers[&key].clone();
compiled_layers.push(Some(precompiled));
self.layers_compiled_per_call.fetch_add(1, Ordering::SeqCst);
}
Ok(compiled_layers)
}
}
}