use std::os::unix::io::OwnedFd;
use std::sync::Arc;
use anyhow::{Context, Result};
use base64::Engine;
use composefs::{
fsverity::FsVerityHashValue,
repository::{ImportContext, Repository},
};
use cstorage::{
CstorLayerService, Image, Layer, Storage, StorageProxy, can_bypass_file_permissions,
spawn_cstor_in_process,
};
use crate::varlink_types::{GetLayerParams, OciProxy as _, StorageLocator};
pub use cstorage::init_if_helper;
use crate::progress::{ComponentId, ProgressEvent, ProgressUnit, SharedReporter};
use crate::{ContentAndVerity, ImportStats, OciDigest, layer_identifier};
type CstorImportResult<ObjectID> = (ContentAndVerity<ObjectID>, ContentAndVerity<ObjectID>);
pub async fn import_from_containers_storage<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
image_id: &str,
reference: Option<&str>,
zerocopy: bool,
storage_root: Option<&std::path::Path>,
additional_image_stores: &[&std::path::Path],
reporter: SharedReporter,
) -> Result<(CstorImportResult<ObjectID>, ImportStats)> {
if can_bypass_file_permissions() {
let storage_root = storage_root.map(|p| p.to_path_buf());
let additional_image_stores: Vec<std::path::PathBuf> = additional_image_stores
.iter()
.map(|p| p.to_path_buf())
.collect();
import_from_containers_storage_direct(
repo,
image_id,
reference,
zerocopy,
storage_root.as_deref(),
&additional_image_stores,
reporter,
)
.await
} else {
if storage_root.is_some() || !additional_image_stores.is_empty() {
anyhow::bail!(
"storage_root and additional_image_stores are not supported in rootless mode"
);
}
import_from_containers_storage_proxied(repo, image_id, reference, zerocopy, reporter).await
}
}
struct ResolvedImageLayers {
image: Image,
layers: Vec<(String, String, OciDigest)>,
}
fn resolve_image_layers(
image_id: &str,
storage_root: Option<std::path::PathBuf>,
additional_image_stores: Vec<std::path::PathBuf>,
) -> Result<ResolvedImageLayers> {
let mut store_paths: Vec<String> = match &storage_root {
Some(root) => vec![root.to_string_lossy().into_owned()],
None => storage_search_paths(),
};
for p in &additional_image_stores {
store_paths.push(p.to_string_lossy().into_owned());
}
let mut stores: Vec<(String, Storage)> = Vec::with_capacity(store_paths.len());
let mut open_errors: Vec<String> = Vec::new();
for path in &store_paths {
match Storage::open(path) {
Ok(s) => stores.push((path.clone(), s)),
Err(e) => open_errors.push(format!("{path}: {e:#}")),
}
}
if stores.is_empty() {
anyhow::bail!(
"Could not open any containers-storage root: {}",
open_errors.join("; ")
);
}
let image = stores
.iter()
.find_map(|(_path, s)| {
Image::open(s, image_id)
.or_else(|_| s.find_image_by_name(image_id))
.ok()
})
.with_context(|| format!("Failed to find image {image_id} in any storage"))?;
let storage_vec: Vec<Storage> = stores
.iter()
.map(|(path, _s)| Storage::open(path))
.collect::<std::result::Result<Vec<_>, _>>()
.context("Failed to re-open stores for storage_layer_ids")?;
let storage_layer_ids = image
.storage_layer_ids(&storage_vec)
.context("Failed to get storage layer IDs from image")?;
let config = image.config().context("Failed to read image config")?;
let diff_ids: Vec<OciDigest> = config
.rootfs()
.diff_ids()
.iter()
.map(|s| s.parse::<OciDigest>().context("parsing diff_id"))
.collect::<Result<_>>()?;
anyhow::ensure!(
storage_layer_ids.len() == diff_ids.len(),
"Layer count mismatch: {} layers in storage, {} diff_ids in config",
storage_layer_ids.len(),
diff_ids.len()
);
let mut layers = Vec::with_capacity(storage_layer_ids.len());
for (storage_layer_id, diff_id) in storage_layer_ids.into_iter().zip(diff_ids) {
let store_path = stores
.iter()
.find_map(|(path, s)| Layer::open(s, &storage_layer_id).ok().map(|_| path.clone()))
.with_context(|| format!("Could not find store containing layer {storage_layer_id}"))?;
layers.push((store_path, storage_layer_id, diff_id));
}
Ok(ResolvedImageLayers { image, layers })
}
fn storage_search_paths() -> Vec<String> {
let mut paths = Vec::new();
if let Ok(root) = std::env::var("CONTAINERS_STORAGE_ROOT") {
paths.push(root);
}
if let Ok(home) = std::env::var("HOME") {
if let Ok(xdg) = std::env::var("XDG_DATA_HOME") {
paths.push(format!("{xdg}/containers/storage"));
}
paths.push(format!("{home}/.local/share/containers/storage"));
}
paths.push("/var/lib/containers/storage".to_string());
if let Ok(opts) = std::env::var("STORAGE_OPTS") {
for item in opts.split(',') {
let item = item.trim();
if let Some(p) = item.strip_prefix("additionalimagestore=") {
paths.push(p.to_string());
}
}
}
paths
}
async fn import_from_containers_storage_direct<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
image_id: &str,
reference: Option<&str>,
zerocopy: bool,
storage_root: Option<&std::path::Path>,
additional_image_stores: &[std::path::PathBuf],
reporter: SharedReporter,
) -> Result<(CstorImportResult<ObjectID>, ImportStats)> {
let mut stats = ImportStats::default();
let mut ctx = ImportContext::default();
let image_id_owned = image_id.to_owned();
let storage_root_owned = storage_root.map(|p| p.to_path_buf());
let additional_owned: Vec<std::path::PathBuf> = additional_image_stores.to_vec();
let resolved = tokio::task::spawn_blocking(move || {
resolve_image_layers(&image_id_owned, storage_root_owned, additional_owned)
})
.await
.context("spawn_blocking(resolve_image_layers) failed")??;
stats.layers = resolved.layers.len() as u64;
let (mut client, server) =
spawn_cstor_in_process(CstorLayerService).context("Failed to spawn CstorLayerService")?;
let mut layer_refs = Vec::with_capacity(resolved.layers.len());
for (store_path, storage_layer_id, diff_id) in &resolved.layers {
let content_id = layer_identifier(diff_id);
let id = ComponentId::from(diff_id.to_string());
let layer_verity = if let Some(existing) = repo.has_stream(&content_id)? {
reporter.report(ProgressEvent::Skipped { id });
stats.layers_already_present += 1;
existing
} else {
reporter.report(ProgressEvent::Started {
id: id.clone(),
total: None,
unit: ProgressUnit::Bytes,
});
let (verity, layer_stats) = import_layer_via_transfer(
repo,
&mut client,
store_path,
storage_layer_id,
diff_id,
zerocopy,
true, &mut ctx,
)
.await?;
let bytes = layer_stats.new_bytes();
stats.merge(&layer_stats);
reporter.report(ProgressEvent::Done {
id,
transferred: bytes,
});
verity
};
layer_refs.push((diff_id.clone(), layer_verity));
}
drop(client);
server.shutdown().await;
reporter.report(ProgressEvent::Message("Layers imported".to_string()));
let repo2 = Arc::clone(repo);
let image = resolved.image;
let reference_owned = reference.map(|s| s.to_owned());
let reporter2 = reporter.clone();
tokio::task::spawn_blocking(move || {
finalize_import(
&repo2,
&image,
&layer_refs,
reference_owned.as_deref(),
&reporter2,
stats,
)
})
.await
.context("spawn_blocking(finalize_import) failed")?
}
#[allow(clippy::too_many_arguments)]
async fn import_layer_via_transfer<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
client: &mut zlink::tokio::unix::Connection,
storage_path: &str,
storage_layer_id: &str,
diff_id: &OciDigest,
zerocopy: bool,
consumer_has_cap_dac_override: bool,
ctx: &mut ImportContext,
) -> Result<(ObjectID, ImportStats)> {
let params = GetLayerParams {
diff_id: None,
storage: Some(StorageLocator {
storage_path: storage_path.to_owned(),
layer_id: storage_layer_id.to_owned(),
}),
consumer_has_cap_dac_override,
};
let stream = client
.get_layer(0, params)
.await
.with_context(|| format!("Oci.GetLayer RPC failed for {storage_layer_id}"))?;
let mut fds: Vec<OwnedFd> = Vec::new();
let mut reply_opt: Option<crate::varlink_types::GetLayerReply> = None;
{
use zlink::futures_util::StreamExt as _;
let mut stream = std::pin::pin!(stream);
while let Some(item) = stream.next().await {
let (result, frame_fds) =
item.with_context(|| format!("Oci.GetLayer stream error for {storage_layer_id}"))?;
let frame_reply = result.map_err(|e| anyhow::anyhow!("Oci.GetLayer error: {e:?}"))?;
reply_opt = Some(frame_reply);
fds.extend(frame_fds);
}
}
let reply = reply_opt.ok_or_else(|| anyhow::anyhow!("Oci.GetLayer yielded no frames"))?;
anyhow::ensure!(
fds.len() >= reply.dir_count as usize + 2,
"Oci.GetLayer: expected at least {} fds (1 pipe + {} dirfd slots + 1 keepalive), got {}",
2 + reply.dir_count,
reply.dir_count,
fds.len()
);
let mut it = fds.into_iter();
let pipe_read: OwnedFd = it.next().expect("checked above: fds.len() >= 1");
let dir_fds: Vec<OwnedFd> = it.by_ref().take(reply.dir_count as usize).collect();
let lifetime_fds: Vec<OwnedFd> = it.collect();
let repo_clone = Arc::clone(repo);
let diff_id_owned = diff_id.clone();
let ctx_moved = std::mem::take(ctx);
let (verity, layer_stats, ctx_returned) = tokio::task::spawn_blocking(move || {
let _lifetime_fds = lifetime_fds;
crate::layer_sync::drain_splitdirfdstream(
repo_clone,
pipe_read,
dir_fds,
&diff_id_owned,
zerocopy,
ctx_moved,
)
})
.await
.context("spawn_blocking(drain_splitdirfdstream) failed")??;
*ctx = ctx_returned;
Ok((verity, layer_stats))
}
async fn import_from_containers_storage_proxied<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
image_id: &str,
reference: Option<&str>,
zerocopy: bool,
reporter: SharedReporter,
) -> Result<(CstorImportResult<ObjectID>, ImportStats)> {
let mut stats = ImportStats::default();
let mut ctx = ImportContext::default();
let image_id_owned = image_id.to_owned();
let resolved =
tokio::task::spawn_blocking(move || resolve_image_layers(&image_id_owned, None, vec![]))
.await
.context("spawn_blocking(resolve_image_layers) failed")??;
stats.layers = resolved.layers.len() as u64;
let mut proxy = StorageProxy::spawn()
.await
.context("spawn userns helper")?
.context("expected helper but got None (can_bypass_file_permissions returned true?)")?;
let mut layer_refs = Vec::with_capacity(resolved.layers.len());
for (store_path, storage_layer_id, diff_id) in &resolved.layers {
let content_id = layer_identifier(diff_id);
let id = ComponentId::from(diff_id.to_string());
let layer_verity = if let Some(existing) = repo.has_stream(&content_id)? {
reporter.report(ProgressEvent::Skipped { id });
stats.layers_already_present += 1;
existing
} else {
reporter.report(ProgressEvent::Started {
id: id.clone(),
total: None,
unit: ProgressUnit::Bytes,
});
let (verity, layer_stats) = import_layer_via_transfer(
repo,
proxy.connection(),
store_path,
storage_layer_id,
diff_id,
zerocopy,
false, &mut ctx,
)
.await?;
let bytes = layer_stats.new_bytes();
stats.merge(&layer_stats);
reporter.report(ProgressEvent::Done {
id,
transferred: bytes,
});
verity
};
layer_refs.push((diff_id.clone(), layer_verity));
}
proxy.shutdown().await.context("shutdown userns helper")?;
reporter.report(ProgressEvent::Message("Layers imported".to_string()));
let repo2 = Arc::clone(repo);
let image = resolved.image;
let reference_owned = reference.map(|s| s.to_owned());
let reporter2 = reporter.clone();
tokio::task::spawn_blocking(move || {
finalize_import(
&repo2,
&image,
&layer_refs,
reference_owned.as_deref(),
&reporter2,
stats,
)
})
.await
.context("spawn_blocking(finalize_import) failed")?
}
fn finalize_import<ObjectID: FsVerityHashValue>(
repo: &Arc<Repository<ObjectID>>,
image: &Image,
layer_refs: &[(OciDigest, ObjectID)],
reference: Option<&str>,
reporter: &SharedReporter,
stats: ImportStats,
) -> Result<(CstorImportResult<ObjectID>, ImportStats)> {
let config_key = format!("sha256:{}", image.id());
let encoded_key = base64::engine::general_purpose::STANDARD.encode(config_key.as_bytes());
let config_json = image
.read_metadata(&encoded_key)
.context("Failed to read config bytes")?;
reporter.report(ProgressEvent::Message(format!(
"Read config ({} bytes)",
config_json.len()
)));
let manifest_json = image
.read_manifest_raw()
.context("Failed to read manifest bytes")?;
reporter.report(ProgressEvent::Message(format!(
"Read manifest ({} bytes)",
manifest_json.len()
)));
let result = crate::layer_sync::finalize_oci_image(
repo,
&manifest_json,
&config_json,
layer_refs,
reference,
)
.context("finalize_oci_image")?;
Ok((result, stats))
}
pub fn parse_containers_storage_ref(imgref: &str) -> Option<&str> {
imgref.strip_prefix("containers-storage:")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_parse_containers_storage_ref() {
assert_eq!(
parse_containers_storage_ref("containers-storage:sha256:abc123"),
Some("sha256:abc123")
);
assert_eq!(
parse_containers_storage_ref("containers-storage:quay.io/fedora:latest"),
Some("quay.io/fedora:latest")
);
assert_eq!(
parse_containers_storage_ref("docker://quay.io/fedora:latest"),
None
);
assert_eq!(parse_containers_storage_ref("sha256:abc123"), None);
}
}