use std::os::unix::fs::FileExt;
use std::os::unix::io::OwnedFd;
use std::sync::Arc;
use anyhow::{Context, Result};
use cap_std_ext::cap_std;
use composefs::digest::{Digest as _, Sha256};
use composefs_splitdirfdstream::{Chunk, SplitdirfdstreamReader, SplitdirfdstreamWriter};
use rustix::fs::{MemfdFlags, fstat, memfd_create};
use composefs::{
INLINE_CONTENT_MAX_V0,
fsverity::FsVerityHashValue,
repository::{ImportContext, ObjectStoreMethod, Repository},
splitstream::{SplitStreamData, SplitStreamWriter},
};
use crate::skopeo::TAR_LAYER_CONTENT_TYPE;
use crate::{ImportStats, OciDigest, layer_identifier, sha256_output_to_digest};
#[derive(Debug, thiserror::Error)]
pub enum VerifiedDrainError {
#[error("layer content does not match declared diff_id: expected {expected}, got {actual}")]
DiffIdMismatch {
expected: String,
actual: String,
},
#[error(transparent)]
Other(#[from] anyhow::Error),
}
fn assert_is_dir(fd: &impl rustix::fd::AsFd, slot: u32, name: &str) -> anyhow::Result<()> {
use rustix::fs::FileType;
let st = fstat(fd).with_context(|| format!("fstat overlay_dir[{slot}]"))?;
let ft = FileType::from_raw_mode(st.st_mode);
if ft != FileType::Directory {
anyhow::bail!(
"overlay_dir[{slot}] is not a directory (file type: {ft:?}, \
mode={:#o}, size={}) — cannot open {name:?}; \
this is likely a bug where a dummy fd (e.g. /dev/null or memfd) \
was stored at a slot that should hold a real diff-dir fd",
st.st_mode,
st.st_size
);
}
Ok(())
}
pub fn drain_splitdirfdstream<ObjectID: FsVerityHashValue>(
repo: Arc<Repository<ObjectID>>,
pipe_read: OwnedFd,
dir_fds: Vec<OwnedFd>,
diff_id: &OciDigest,
zerocopy: bool,
mut ctx: ImportContext,
) -> Result<(ObjectID, ImportStats, ImportContext)> {
let overlay_dirs: Vec<cap_std::fs::Dir> =
dir_fds.into_iter().map(cap_std::fs::Dir::from).collect();
let mut writer = repo.create_stream(TAR_LAYER_CONTENT_TYPE)?;
let content_id = layer_identifier(diff_id);
let mut reader = SplitdirfdstreamReader::new(std::fs::File::from(pipe_read));
let mut inline_buf = Vec::new();
let mut stats = ImportStats::default();
drain_splitdirfdstream_inner(
&repo,
&mut writer,
&mut reader,
&overlay_dirs,
zerocopy,
&mut stats,
&mut ctx,
&mut inline_buf,
None,
)?;
let verity = repo.write_stream(writer, &content_id, None)?;
Ok((verity, stats, ctx))
}
#[allow(clippy::too_many_arguments)]
fn drain_splitdirfdstream_inner<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
writer: &mut SplitStreamWriter<ObjectID>,
reader: &mut SplitdirfdstreamReader<std::fs::File>,
overlay_dirs: &[cap_std::fs::Dir],
zerocopy: bool,
stats: &mut ImportStats,
ctx: &mut ImportContext,
inline_buf: &mut Vec<u8>,
mut hasher: Option<&mut Sha256>,
) -> Result<()> {
while let Some(chunk) = reader.next_chunk().context("splitdirfdstream read error")? {
match chunk {
Chunk::Metadata(data) => {
if let Some(ref mut h) = hasher {
h.update(data);
}
stats.bytes_inlined += data.len() as u64;
writer.write_inline(data);
}
Chunk::InlineData(data) => {
if let Some(ref mut h) = hasher {
h.update(data);
}
let length = data.len() as u64;
if should_inline(length) {
stats.bytes_inlined += length;
writer.write_inline(data);
} else {
process_file_content(
repo,
writer,
stats,
ctx,
file_content_to_memfd(data)?,
length,
"<file-content>",
false,
inline_buf,
)?;
}
}
Chunk::FileBackedData {
dirfd_index,
length,
filename,
} => {
let name = std::str::from_utf8(filename).with_context(|| {
format!("non-utf8 filename in splitdirfdstream: {filename:?}")
})?;
let dir = overlay_dirs.get(dirfd_index as usize).with_context(|| {
format!(
"dirfd_index {dirfd_index} out of range (dir_fds.len={})",
overlay_dirs.len()
)
})?;
assert_is_dir(dir, dirfd_index, name)?;
let fd = dir
.open(name)
.map(OwnedFd::from)
.with_context(|| format!("open {name:?} in overlay dir[{dirfd_index}]"))?;
use rustix::fs::FileType;
let st = fstat(&fd).with_context(|| format!("fstat {name:?}"))?;
let ft = FileType::from_raw_mode(st.st_mode);
if ft != FileType::RegularFile {
anyhow::bail!(
"object {name:?} in overlay dir[{dirfd_index}] is not a regular file (file type: {ft:?}, mode={:#o})",
st.st_mode
);
}
let fd_to_process = if let Some(ref mut h) = hasher {
let obj_file = std::fs::File::from(fd);
let actual_size = st.st_size as u64;
if actual_size != length {
anyhow::bail!(
"object {name}: declared length {length} != actual size {actual_size}"
);
}
hash_fd_contents(&obj_file, length, h)
.with_context(|| format!("hashing object {name}"))?;
obj_file.into()
} else {
fd
};
process_file_content(
repo,
writer,
stats,
ctx,
fd_to_process,
length,
name,
zerocopy,
inline_buf,
)?;
}
}
}
Ok(())
}
fn file_content_to_memfd(data: &[u8]) -> Result<OwnedFd> {
let memfd = memfd_create(c"composefs-filecontent", MemfdFlags::CLOEXEC)
.context("memfd_create for FileContent chunk")?;
rustix::io::write(&memfd, data).context("writing FileContent to memfd")?;
Ok(memfd)
}
fn hash_fd_contents(fd: &std::fs::File, len: u64, hasher: &mut Sha256) -> Result<()> {
const BUF_SIZE: usize = 65536;
let mut buf = [0u8; BUF_SIZE];
let mut remaining = len;
let mut offset = 0u64;
while remaining > 0 {
let to_read = remaining.min(BUF_SIZE as u64) as usize;
let n = fd
.read_at(&mut buf[..to_read], offset)
.context("read_at while hashing fd contents")?;
if n == 0 {
anyhow::bail!(
"unexpected EOF at offset {offset} hashing fd (expected {len} bytes total)"
);
}
hasher.update(&buf[..n]);
offset += n as u64;
remaining -= n as u64;
}
Ok(())
}
pub fn drain_splitdirfdstream_verified<ObjectID: FsVerityHashValue>(
repo: Arc<Repository<ObjectID>>,
pipe_read: OwnedFd,
dir_fds: Vec<OwnedFd>,
diff_id: &OciDigest,
zerocopy: bool,
mut ctx: ImportContext,
) -> Result<(ObjectID, ImportStats, ImportContext), VerifiedDrainError> {
let algorithm = diff_id.algorithm().as_ref();
if algorithm != "sha256" {
return Err(VerifiedDrainError::Other(anyhow::anyhow!(
"unsupported diff_id algorithm {algorithm:?}: only sha256 is supported"
)));
}
let overlay_dirs: Vec<cap_std::fs::Dir> =
dir_fds.into_iter().map(cap_std::fs::Dir::from).collect();
let mut writer = repo
.create_stream(TAR_LAYER_CONTENT_TYPE)
.context("create_stream")?;
let content_id = layer_identifier(diff_id);
let mut reader = SplitdirfdstreamReader::new(std::fs::File::from(pipe_read));
let mut inline_buf = Vec::new();
let mut stats = ImportStats::default();
let mut hasher = Sha256::new();
drain_splitdirfdstream_inner(
&repo,
&mut writer,
&mut reader,
&overlay_dirs,
zerocopy,
&mut stats,
&mut ctx,
&mut inline_buf,
Some(&mut hasher),
)?;
let actual_digest = sha256_output_to_digest(hasher.finalize());
let actual_str = actual_digest.to_string();
let expected_str = diff_id.to_string();
if actual_str != expected_str {
return Err(VerifiedDrainError::DiffIdMismatch {
expected: expected_str,
actual: actual_str,
});
}
let verity = repo
.write_stream(writer, &content_id, None)
.context("write_stream")?;
Ok((verity, stats, ctx))
}
pub(crate) fn should_inline(size: u64) -> bool {
(size as usize) <= INLINE_CONTENT_MAX_V0
}
#[allow(clippy::too_many_arguments)]
pub fn process_file_content<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
writer: &mut SplitStreamWriter<ObjectID>,
stats: &mut ImportStats,
ctx: &mut ImportContext,
fd: OwnedFd,
size: u64,
name: &str,
zerocopy: bool,
inline_buf: &mut Vec<u8>,
) -> Result<()> {
let file = std::fs::File::from(fd);
if !should_inline(size) {
let (object_id, method) = if zerocopy {
repo.ensure_object_from_file_zerocopy(&file, size, ctx)
} else {
repo.ensure_object_from_file(&file, size, ctx)
}
.with_context(|| format!("Failed to store object for {}", name))?;
match method {
ObjectStoreMethod::Reflinked => {
stats.objects_reflinked += 1;
stats.bytes_reflinked += size;
}
ObjectStoreMethod::Hardlinked => {
stats.objects_hardlinked += 1;
stats.bytes_hardlinked += size;
}
ObjectStoreMethod::Copied => {
stats.objects_copied += 1;
stats.bytes_copied += size;
}
ObjectStoreMethod::AlreadyPresent => {
stats.objects_already_present += 1;
}
}
writer.add_external_size(size);
writer.write_reference(object_id)?;
} else {
inline_buf.resize(size as usize, 0);
file.read_exact_at(inline_buf, 0)?;
stats.bytes_inlined += size;
writer.write_inline(inline_buf);
}
Ok(())
}
fn object_size<ObjectID: FsVerityHashValue>(
repo: &Repository<ObjectID>,
id: &ObjectID,
) -> Result<u64> {
let fd = repo
.open_object(id)
.context("Opening object for size query")?;
let stat = fstat(&fd).context("fstat on object fd")?;
Ok(stat.st_size as u64)
}
pub fn produce_layer_splitdirfdstream<ObjectID: FsVerityHashValue, W: std::io::Write>(
repo: &Repository<ObjectID>,
layer_verity: &ObjectID,
objects_dirfd_index: u32,
out: W,
) -> Result<()> {
let mut reader = repo
.open_stream("", Some(layer_verity), None)
.context("Opening layer splitstream")?;
let mut writer = SplitdirfdstreamWriter::new(out);
reader
.for_each_chunk(|chunk| {
match chunk {
SplitStreamData::Inline(data) => {
writer.write_metadata(&data).map_err(anyhow::Error::from)?;
}
SplitStreamData::External(id) => {
let size = object_size(repo, &id)?;
let pathname = id.to_object_pathname();
writer
.write_file_backed_data(objects_dirfd_index, size, pathname.as_bytes())
.map_err(anyhow::Error::from)?;
}
}
Ok(())
})
.context("Walking layer splitstream chunks")?;
writer.finish().map_err(anyhow::Error::from)?;
Ok(())
}
pub type FinalizeResult<ObjectID> = (
crate::ContentAndVerity<ObjectID>,
crate::ContentAndVerity<ObjectID>,
);
pub fn finalize_oci_image<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
manifest_json: &[u8],
config_json: &[u8],
layer_refs: &[(OciDigest, ObjectID)],
name: Option<&str>,
) -> anyhow::Result<FinalizeResult<ObjectID>> {
use crate::oci_image::manifest_identifier;
use crate::skopeo::{OCI_CONFIG_CONTENT_TYPE, OCI_MANIFEST_CONTENT_TYPE};
use crate::{config_identifier, sha256_content_digest};
let config_digest = sha256_content_digest(config_json);
let content_id = config_identifier(&config_digest);
let config_verity = if let Some(existing) = repo.has_stream(&content_id)? {
existing
} else {
let mut writer = repo.create_stream(OCI_CONFIG_CONTENT_TYPE)?;
for (diff_id, verity) in layer_refs {
let key: &str = diff_id.as_ref();
writer.add_named_stream_ref(key, verity);
}
writer.write_external(config_json)?;
repo.write_stream(writer, &content_id, None)?
};
let manifest_digest = sha256_content_digest(manifest_json);
let manifest_content_id = manifest_identifier(&manifest_digest);
let manifest_verity = if let Some(existing) = repo.has_stream(&manifest_content_id)? {
existing
} else {
let mut writer = repo.create_stream(OCI_MANIFEST_CONTENT_TYPE)?;
let config_ref_key = format!("config:{config_digest}");
writer.add_named_stream_ref(&config_ref_key, &config_verity);
for (diff_id, verity) in layer_refs {
let key: &str = diff_id.as_ref();
writer.add_named_stream_ref(key, verity);
}
writer.write_external(manifest_json)?;
repo.write_stream(writer, &manifest_content_id, None)?
};
let existing_erofs = crate::composefs_erofs_for_manifest(
repo,
&manifest_digest,
Some(&manifest_verity),
repo.erofs_version(),
)?;
if existing_erofs.is_none() {
let erofs = crate::ensure_oci_composefs_erofs(
repo,
&manifest_digest,
Some(&manifest_verity),
name,
)?;
if erofs.is_none() {
if let Some(n) = name {
crate::oci_image::tag_image(repo, &manifest_digest, n)?;
}
}
} else if let Some(n) = name {
crate::oci_image::tag_image(repo, &manifest_digest, n)?;
}
let config_verity = repo
.has_stream(&content_id)?
.context("config splitstream missing after finalization")?;
let manifest_verity = repo
.has_stream(&manifest_content_id)?
.context("manifest splitstream missing after finalization")?;
Ok((
(manifest_digest, manifest_verity),
(config_digest, config_verity),
))
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write as _;
use composefs::fsverity::Sha256HashValue;
use composefs::repository::RepositoryConfig;
use composefs_splitdirfdstream::reconstruct;
#[test]
fn test_should_inline_boundary() {
assert!(should_inline(0), "size 0 should be inlined");
assert!(should_inline(1), "size 1 should be inlined");
assert!(
should_inline(INLINE_CONTENT_MAX_V0 as u64),
"size == INLINE_CONTENT_MAX_V0 ({INLINE_CONTENT_MAX_V0}) should be inlined"
);
assert!(
!should_inline(INLINE_CONTENT_MAX_V0 as u64 + 1),
"size INLINE_CONTENT_MAX_V0+1 ({}) should NOT be inlined",
INLINE_CONTENT_MAX_V0 + 1
);
for size in [128u64, 4096, 65536, 1024 * 1024] {
assert!(
!should_inline(size),
"size {size} should NOT be inlined (well above threshold)"
);
}
}
fn create_test_repo() -> (Arc<Repository<Sha256HashValue>>, tempfile::TempDir) {
let tempdir = tempfile::TempDir::new().unwrap();
let (repo, _) = Repository::init_path(
rustix::fs::CWD,
&tempdir.path().join("repo"),
RepositoryConfig::default().set_insecure(),
)
.unwrap();
(Arc::new(repo), tempdir)
}
fn tmpfile_of(len: usize) -> OwnedFd {
let mut f = tempfile::tempfile().unwrap();
let data: Vec<u8> = (0..len).map(|i| (i % 251) as u8).collect();
f.write_all(&data).unwrap();
f.into()
}
#[test]
fn test_process_file_content_inline_vs_external() {
let (repo, _tempdir) = create_test_repo();
let cases = [
(0usize, false),
(1, false),
(INLINE_CONTENT_MAX_V0, false),
(INLINE_CONTENT_MAX_V0 + 1, true),
(4096, true),
(256 * 1024, true),
];
for (size, expect_external) in cases {
let mut writer = repo.create_stream(TAR_LAYER_CONTENT_TYPE).unwrap();
let mut stats = ImportStats::default();
let mut ctx = ImportContext::default();
let mut inline_buf = Vec::new();
let before_inlined = stats.bytes_inlined;
process_file_content(
&repo,
&mut writer,
&mut stats,
&mut ctx,
tmpfile_of(size),
size as u64,
"test-file",
false,
&mut inline_buf,
)
.unwrap();
let objects_written = stats.objects_reflinked
+ stats.objects_hardlinked
+ stats.objects_copied
+ stats.objects_already_present;
if expect_external {
assert_eq!(
stats.bytes_inlined, before_inlined,
"size {size}: external file must not change bytes_inlined"
);
assert_eq!(
objects_written, 1,
"size {size}: exactly one external object expected"
);
} else {
assert_eq!(
stats.bytes_inlined,
before_inlined + size as u64,
"size {size}: inline file must add its bytes to bytes_inlined"
);
assert_eq!(
objects_written, 0,
"size {size}: inline file must not write an object"
);
}
}
}
#[test]
fn test_memfd_file_content_uses_copy_fallback() {
let (repo, _tempdir) = create_test_repo();
let size: usize = INLINE_CONTENT_MAX_V0 + 1; let data: Vec<u8> = (0..size).map(|i| (i % 251) as u8).collect();
let memfd = file_content_to_memfd(&data).unwrap();
let mut writer = repo.create_stream(TAR_LAYER_CONTENT_TYPE).unwrap();
let mut stats = ImportStats::default();
let mut ctx = ImportContext::default();
let mut inline_buf = Vec::new();
process_file_content(
&repo,
&mut writer,
&mut stats,
&mut ctx,
memfd,
size as u64,
"memfd-test",
false,
&mut inline_buf,
)
.unwrap();
assert_eq!(stats.objects_copied, 1, "memfd content should be copied");
assert_eq!(stats.bytes_copied, size as u64);
}
fn build_tar_layer(file_sizes: &[usize]) -> Vec<u8> {
let mut builder = ::tar::Builder::new(vec![]);
for &size in file_sizes {
let content: Vec<u8> = (0..size).map(|i| (i % 251) as u8).collect();
let mut header = ::tar::Header::new_ustar();
header.set_uid(0);
header.set_gid(0);
header.set_mode(0o644);
header.set_entry_type(::tar::EntryType::Regular);
header.set_size(size as u64);
builder
.append_data(
&mut header,
format!("file_{size}_{:08x}", size),
&content[..],
)
.unwrap();
}
builder.into_inner().unwrap()
}
async fn assert_produce_eq_cat(
repo: &Arc<Repository<Sha256HashValue>>,
verity: &Sha256HashValue,
) {
use std::os::fd::AsFd as _;
let mut expected = Vec::<u8>::new();
let mut reader = repo
.open_stream("", Some(verity), Some(TAR_LAYER_CONTENT_TYPE))
.expect("open_stream for cat");
reader
.cat(repo, &mut expected)
.expect("cat on layer splitstream");
let mut stream_buf = Vec::<u8>::new();
produce_layer_splitdirfdstream(repo, verity, 0, &mut stream_buf)
.expect("produce_layer_splitdirfdstream");
let objects_dir_fd = repo.objects_dir().expect("objects_dir");
let dirfds = [objects_dir_fd.as_fd()];
let mut actual = Vec::<u8>::new();
reconstruct(stream_buf.as_slice(), &dirfds, &mut actual)
.expect("reconstruct splitdirfdstream");
similar_asserts::assert_eq!(
actual,
expected,
"produce->reconstruct must equal cat for verity={verity:?}"
);
}
#[tokio::test]
async fn test_produce_reconstruct_eq_cat() {
let (repo, _tempdir) = create_test_repo();
let cases: &[(&str, &[usize])] = &[
("empty", &[]),
("all_inline", &[0, 1, 10, 64]),
("single_external", &[65]),
("large_external", &[4096, 200_000]),
("mixed", &[0, 10, 64, 65, 4096, 200_000]),
];
for (label, sizes) in cases {
let tar_bytes = build_tar_layer(sizes);
let diff_id = crate::sha256_content_digest(&tar_bytes);
let (verity, _stats) = crate::import_layer(&repo, &diff_id, None, tar_bytes.as_slice())
.await
.unwrap_or_else(|e| panic!("import_layer failed for {label}: {e}"));
assert_produce_eq_cat(&repo, &verity).await;
}
}
fn produce_to_pipe(
repo_a: Arc<Repository<Sha256HashValue>>,
verity: Sha256HashValue,
) -> (
OwnedFd,
Vec<OwnedFd>,
tokio::task::JoinHandle<anyhow::Result<()>>,
) {
use std::os::fd::AsFd as _;
let (pipe_read, pipe_write) =
rustix::pipe::pipe_with(rustix::pipe::PipeFlags::CLOEXEC).expect("pipe");
let objects_dir = repo_a.objects_dir().expect("objects_dir");
let objects_owned = rustix::io::dup(objects_dir.as_fd()).expect("dup objects_dir");
let handle = tokio::task::spawn_blocking(move || {
let wf = std::fs::File::from(pipe_write);
produce_layer_splitdirfdstream(&repo_a, &verity, 0, wf)
});
(pipe_read, vec![objects_owned], handle)
}
async fn join_producer(handle: tokio::task::JoinHandle<anyhow::Result<()>>) {
handle
.await
.expect("producer task panicked")
.expect("producer must succeed");
}
async fn drain_producer(handle: tokio::task::JoinHandle<anyhow::Result<()>>) {
let _ = handle.await.expect("producer task panicked");
}
#[tokio::test]
async fn test_verified_drain_correct_diff_id() {
let (repo_a, _td_a) = create_test_repo();
let (repo_b, _td_b) = create_test_repo();
let tar_bytes = build_tar_layer(&[10, 128 * 1024]); let diff_id = crate::sha256_content_digest(&tar_bytes);
let (verity_a, _) = crate::import_layer(&repo_a, &diff_id, None, tar_bytes.as_slice())
.await
.expect("import_layer into repo_a");
let (pipe_read, dir_fds, producer) = produce_to_pipe(repo_a, verity_a.clone());
let repo_b_clone = repo_b.clone();
let diff_id_clone = diff_id.clone();
let result = tokio::task::spawn_blocking(move || {
drain_splitdirfdstream_verified(
repo_b_clone,
pipe_read,
dir_fds,
&diff_id_clone,
false,
composefs::repository::ImportContext::default(),
)
})
.await
.expect("spawn_blocking");
let (verity_b, _stats, _ctx) = result.expect("verified drain must succeed");
join_producer(producer).await;
let content_id = crate::layer_content_id(&diff_id);
assert!(
repo_b
.has_stream(&content_id)
.expect("has_stream")
.is_some(),
"repo_b must have the layer stream after verified drain"
);
assert_eq!(
verity_a, verity_b,
"verity hash must be identical across repos"
);
}
#[tokio::test]
async fn test_verified_drain_wrong_diff_id() {
let (repo_a, _td_a) = create_test_repo();
let (repo_b, _td_b) = create_test_repo();
let tar_bytes = build_tar_layer(&[10, 128 * 1024]);
let correct_diff_id = crate::sha256_content_digest(&tar_bytes);
let (verity_a, _) =
crate::import_layer(&repo_a, &correct_diff_id, None, tar_bytes.as_slice())
.await
.expect("import_layer into repo_a");
let wrong_diff_id: crate::OciDigest =
"sha256:0000000000000000000000000000000000000000000000000000000000000000"
.parse()
.unwrap();
let (pipe_read, dir_fds, producer) = produce_to_pipe(repo_a, verity_a);
let repo_b_clone = repo_b.clone();
let wrong_diff_id_clone = wrong_diff_id.clone();
let result = tokio::task::spawn_blocking(move || {
drain_splitdirfdstream_verified(
repo_b_clone,
pipe_read,
dir_fds,
&wrong_diff_id_clone,
false,
composefs::repository::ImportContext::default(),
)
})
.await
.expect("spawn_blocking");
join_producer(producer).await;
match result {
Err(VerifiedDrainError::DiffIdMismatch { expected, actual }) => {
assert_eq!(
expected,
wrong_diff_id.to_string(),
"expected field must be the wrong diff_id"
);
assert_eq!(
actual,
correct_diff_id.to_string(),
"actual field must be the real content hash"
);
}
other => panic!("expected DiffIdMismatch, got {other:?}"),
}
let wrong_content_id = crate::layer_content_id(&wrong_diff_id);
assert!(
repo_b
.has_stream(&wrong_content_id)
.expect("has_stream")
.is_none(),
"repo_b must NOT have a committed stream for the wrong diff_id"
);
}
#[tokio::test]
async fn test_verified_drain_rejects_non_sha256() {
let (repo_a, _td_a) = create_test_repo();
let (repo_b, _td_b) = create_test_repo();
let tar_bytes = build_tar_layer(&[10, 128 * 1024]);
let diff_id = crate::sha256_content_digest(&tar_bytes);
let (verity_a, _) = crate::import_layer(&repo_a, &diff_id, None, tar_bytes.as_slice())
.await
.expect("import_layer into repo_a");
let sha512_diff_id: crate::OciDigest = format!("sha512:{}", "0".repeat(128))
.parse()
.expect("valid sha512 digest");
let (pipe_read, dir_fds, producer) = produce_to_pipe(repo_a, verity_a);
let repo_b_clone = repo_b.clone();
let result = tokio::task::spawn_blocking(move || {
drain_splitdirfdstream_verified(
repo_b_clone,
pipe_read,
dir_fds,
&sha512_diff_id,
false,
composefs::repository::ImportContext::default(),
)
})
.await
.expect("spawn_blocking");
drain_producer(producer).await;
match result {
Err(VerifiedDrainError::Other(e)) => {
let msg = format!("{e:#}");
assert!(
msg.contains("sha256") && msg.contains("sha512"),
"error should explain the algorithm restriction, got: {msg}"
);
}
other => panic!("expected Other(unsupported algorithm) error, got {other:?}"),
}
}
#[tokio::test]
async fn test_finalize_oci_image() {
let (repo, _tempdir) = create_test_repo();
let tar1 = crate::test_util::build_oci_tar_layer(10);
let tar2 = crate::test_util::build_oci_tar_layer(128 * 1024);
let diff_id1 = crate::sha256_content_digest(&tar1);
let diff_id2 = crate::sha256_content_digest(&tar2);
let (verity1, _) = crate::import_layer(&repo, &diff_id1, None, tar1.as_slice())
.await
.expect("import layer 1");
let (verity2, _) = crate::import_layer(&repo, &diff_id2, None, tar2.as_slice())
.await
.expect("import layer 2");
let diff_ids = vec![diff_id1.to_string(), diff_id2.to_string()];
let config_json = crate::test_util::make_config_json(&diff_ids);
let config_digest = crate::sha256_content_digest(&config_json);
let manifest_json =
crate::test_util::make_manifest_json(&config_json, config_digest.as_ref(), &diff_ids);
let layer_refs = vec![(diff_id1.clone(), verity1), (diff_id2.clone(), verity2)];
let ((manifest_digest, manifest_verity), (out_config_digest, config_verity)) =
finalize_oci_image(
&repo,
&manifest_json,
&config_json,
&layer_refs,
Some("test:v1"),
)
.expect("finalize_oci_image");
assert!(!manifest_digest.to_string().is_empty());
assert!(!out_config_digest.to_string().is_empty());
use crate::oci_image::manifest_identifier;
let manifest_id = manifest_identifier(&manifest_digest);
let config_id = crate::config_identifier(&out_config_digest);
assert!(
repo.has_stream(&manifest_id)
.expect("has_stream manifest")
.is_some(),
"manifest splitstream must exist"
);
assert!(
repo.has_stream(&config_id)
.expect("has_stream config")
.is_some(),
"config splitstream must exist"
);
let stored_manifest_verity = repo
.has_stream(&manifest_id)
.unwrap()
.expect("manifest verity must be stored");
assert_eq!(
manifest_verity, stored_manifest_verity,
"returned manifest_verity must match stored"
);
let stored_config_verity = repo
.has_stream(&config_id)
.unwrap()
.expect("config verity must be stored");
assert_eq!(
config_verity, stored_config_verity,
"returned config_verity must match stored"
);
let erofs = crate::composefs_erofs_for_manifest(
&repo,
&manifest_digest,
Some(&manifest_verity),
repo.erofs_version(),
)
.expect("composefs_erofs_for_manifest");
assert!(
erofs.is_some(),
"EROFS image must exist after finalize_oci_image for a container image"
);
let ((md2, _mv2), (cd2, _cv2)) = finalize_oci_image(
&repo,
&manifest_json,
&config_json,
&layer_refs,
Some("test:v1"),
)
.expect("finalize_oci_image idempotent");
assert_eq!(manifest_digest, md2, "idempotent call: manifest_digest");
assert_eq!(out_config_digest, cd2, "idempotent call: config_digest");
}
}