use std::time::Duration;
use crate::{
ContainerAsync, ContainerRequest, Image,
core::{
client::Client,
containers::{async_container::ContainerLogSource, request::DEFAULT_STARTUP_TIMEOUT},
error::{Result, WaitContainerError},
wait::WaitFor,
},
};
#[cfg(target_os = "macos")]
use crate::core::error::ClientError;
#[cfg(target_os = "macos")]
use crate::core::util::{is_valid_container_id, unique_suffix};
#[cfg(target_os = "macos")]
use crate::core::{
containers::request::ExtraHost,
copy::{CopyDataSource, CopyToContainer},
image::ExecCommand,
wait::CmdWaitFor,
};
#[cfg(target_os = "linux")]
use crate::core::copy::{CopyDataSource, CopyToContainer};
#[expect(async_fn_in_trait)]
pub trait AsyncRunner<I: Image> {
async fn start(self) -> Result<ContainerAsync<I>>;
async fn pull_image(self) -> Result<ContainerRequest<I>>;
}
impl<T, I> AsyncRunner<I> for T
where
T: Into<ContainerRequest<I>> + Send,
I: Image,
{
async fn start(self) -> Result<ContainerAsync<I>> {
#[cfg_attr(target_os = "macos", expect(unused_mut))]
let mut container_req = self.into();
let client = Client::detect()?;
#[cfg(target_os = "macos")]
{
#[expect(clippy::infallible_destructuring_match)]
let client = match client {
Client::MacOs(c) => c,
#[cfg(target_os = "linux")]
Client::Linux(_) => unreachable!("macOS block is not compiled on Linux"),
};
if container_req.health_check().is_some() {
return Err(crate::Error::other(
"with_health_check() is not supported on macOS",
));
}
let id = container_req
.container_name()
.clone()
.unwrap_or_else(|| format!("c-{}", unique_suffix()));
if !is_valid_container_id(&id) {
return Err(crate::Error::other(format!(
"invalid container id {id:?}: must match ^[a-zA-Z0-9][a-zA-Z0-9_.-]+$ and be at most 63 characters"
)));
}
crate::core::client::container_cfg::reject_sctp_ports(&container_req)?;
if let Some(ports) = container_req.ports() {
crate::core::containers::request::reject_duplicate_mapped_ports(ports)?;
}
let descriptor = container_req.descriptor();
let (arch, _rosetta, resolve_platform) =
crate::core::client::container_cfg::normalize_platform(
container_req.platform().as_deref(),
);
let desc_raw =
resolve_or_pull_macos(&client, &descriptor, arch, resolve_platform).await?;
let image_config = crate::core::client::image_config::resolve_image_config(
&client,
&desc_raw,
resolve_platform,
)
.await?;
let kernel = client.get_default_kernel().await?;
let volume_resolutions =
crate::core::client::xpc_client::resolve_volumes(&client, &container_req).await?;
let cfg = crate::core::client::container_cfg::build_config(
&container_req,
&id,
&desc_raw,
&image_config,
&volume_resolutions,
)?;
if let Err(e) = client.create_container(&cfg, kernel).await {
rollback_remove(&id, client.remove(&id, true)).await;
return Err(e);
}
#[cfg(feature = "watchdog")]
if !matches!(
crate::core::env::Config.command(),
crate::core::env::Command::Keep
) {
crate::watchdog::register(&id);
}
if let Err(e) = client.bootstrap_container(&id).await {
rollback_remove(&id, client.remove(&id, true)).await;
return Err(e);
}
if let Err(e) = client.start_process(&id).await {
rollback_remove(&id, client.remove(&id, true)).await;
return Err(e);
}
if let Err(e) = copy_to_sources(&client, &id, &container_req).await {
rollback_remove(&id, client.remove(&id, true)).await;
return Err(e);
}
let startup_timeout = container_req
.startup_timeout()
.unwrap_or(DEFAULT_STARTUP_TIMEOUT);
let ready_conditions = container_req.ready_conditions();
let wait_state = crate::core::containers::async_container::new_wait_state();
{
let generation = wait_state
.lock()
.expect("wait state mutex must not be poisoned while spawning exit code waiter")
.generation();
crate::core::containers::async_container::spawn_exit_code_waiter(
id.clone(),
wait_state.clone(),
generation,
);
}
let hosts: Vec<(String, ExtraHost)> = container_req
.hosts()
.map(|(name, host)| (name.into_owned(), *host))
.collect();
let log_source = match client.logs(&id).await {
Ok((out, err)) => ContainerLogSource::Fd {
stdout: out,
stderr: err,
},
Err(e) => {
if crate::core::containers::async_container::ready_conditions_require_log(
&ready_conditions,
) {
rollback_remove(&id, client.remove(&id, true)).await;
return Err(log_fd_required_error(&e));
}
tracing::warn!("failed to get log fds: {e}");
ContainerLogSource::None
}
};
let container = ContainerAsync::new(
id,
Client::MacOs(client),
container_req,
wait_state,
log_source,
);
if let Err(e) = apply_extra_hosts(&container, &hosts).await {
return cleanup_on_ready_failure(container, e).await;
}
run_ready_sequence(container, startup_timeout, ready_conditions).await
}
#[cfg(target_os = "linux")]
{
#[expect(clippy::infallible_destructuring_match)]
let client = match client {
Client::Linux(c) => c,
#[cfg(target_os = "macos")]
Client::MacOs(_) => unreachable!("Linux block is not compiled on macOS"),
};
if let Some(msg) = linux_unsupported_request_reason(&container_req) {
return Err(crate::Error::other(msg));
}
if let Some(ports) = container_req.ports() {
crate::core::containers::request::reject_duplicate_mapped_ports(ports)?;
}
let descriptor = container_req.descriptor();
let platform = container_req.platform();
let _desc_raw =
resolve_or_pull_linux(&client, &descriptor, platform.as_deref()).await?;
let config = build_container_config(&container_req);
let id = client.create_container(config).await?;
if let Err(e) = copy_to_sources_linux(&client, &id, &container_req).await {
rollback_remove(&id, client.remove(&id, true)).await;
return Err(e);
}
if let Err(e) = client.start_container(&id).await {
rollback_remove(&id, client.remove(&id, true)).await;
return Err(e);
}
let startup_timeout = container_req
.startup_timeout()
.unwrap_or(DEFAULT_STARTUP_TIMEOUT);
let ready_conditions = container_req.ready_conditions();
let log_required = linux_log_stream_required(&container_req, &ready_conditions);
let log_consumers = std::mem::take(&mut container_req.log_consumers);
let (log_source, stored_consumers) =
start_linux_log_stream(&client, &id, log_consumers, log_required).await?;
let wait_state = crate::core::containers::async_container::new_wait_state();
{
let generation = wait_state
.lock()
.expect("wait state mutex must not be poisoned while spawning exit code waiter")
.generation();
crate::core::containers::async_container::spawn_exit_code_waiter(
client.clone(),
id.clone(),
wait_state.clone(),
generation,
);
}
let container = ContainerAsync::new(
id,
Client::Linux(client),
container_req,
wait_state,
log_source,
stored_consumers,
);
run_ready_sequence(container, startup_timeout, ready_conditions).await
}
}
async fn pull_image(self) -> Result<ContainerRequest<I>> {
let container_req = self.into();
let descriptor = container_req.descriptor();
match Client::detect()? {
#[cfg(target_os = "macos")]
Client::MacOs(c) => {
let (arch, _, resolve_platform) =
crate::core::client::container_cfg::normalize_platform(
container_req.platform().as_deref(),
);
c.pull_image(&descriptor, resolve_platform.map(|_| arch))
.await?;
}
#[cfg(target_os = "linux")]
Client::Linux(c) => {
c.pull_image(&descriptor, container_req.platform().as_deref())
.await?;
}
}
Ok(container_req)
}
}
#[cfg(target_os = "macos")]
async fn resolve_or_pull_macos(
client: &crate::core::client::XpcClient,
descriptor: &str,
arch: &str,
resolve_platform: Option<&str>,
) -> Result<String> {
if resolve_platform == Some("linux/amd64") {
client.pull_image(descriptor, Some(arch)).await?;
return client.resolve_image_descriptor(descriptor).await;
}
match client.resolve_image_descriptor(descriptor).await {
Ok(d) => Ok(d),
Err(crate::Error::Client(ClientError::ImageNotFound(_))) => {
client
.pull_image(descriptor, resolve_platform.map(|_| arch))
.await?;
client.resolve_image_descriptor(descriptor).await
}
Err(e) => Err(e),
}
}
#[cfg(target_os = "linux")]
async fn resolve_or_pull_linux(
client: &crate::core::client::DockerClient,
descriptor: &str,
platform: Option<&str>,
) -> Result<String> {
match client.resolve_image_descriptor(descriptor, platform).await {
Ok(d) => Ok(d),
Err(_) => {
client.pull_image(descriptor, platform).await?;
client.resolve_image_descriptor(descriptor, platform).await
}
}
}
async fn run_ready_sequence<I: Image>(
container: ContainerAsync<I>,
startup_timeout: Duration,
ready_conditions: Vec<WaitFor>,
) -> Result<ContainerAsync<I>> {
let result: Result<()> = async {
let state = container.container_state().await?;
for cmd in container.image().exec_before_ready(state)? {
container.exec(cmd).await?;
}
tokio::time::timeout(
startup_timeout,
container.block_until_ready(ready_conditions),
)
.await
.map_err(|_| WaitContainerError::StartupTimeout {
id: container.id().to_string(),
timeout: startup_timeout,
})??;
let state = container.container_state().await?;
for cmd in container.image().exec_after_start(state)? {
container.exec(cmd).await?;
}
Ok(())
}
.await;
if let Err(e) = result {
return cleanup_on_ready_failure(container, e).await;
}
Ok(container)
}
async fn cleanup_on_ready_failure<I: Image>(
container: ContainerAsync<I>,
error: crate::Error,
) -> Result<ContainerAsync<I>> {
let id = container.id().to_string();
if matches!(
crate::core::env::Config.command(),
crate::core::env::Command::Remove
) && let Err(re) = container.rm().await
{
tracing::error!("failed to remove container {id} after startup failure: {re}");
}
Err(error)
}
#[cfg(target_os = "linux")]
fn build_container_config<I: Image>(
req: &ContainerRequest<I>,
) -> crate::core::client::ContainerConfig {
use crate::core::containers::request::PortMapping;
if let Some(ports) = req.ports() {
crate::core::containers::request::reject_duplicate_mapped_ports(ports)
.expect("pull 前検証で重複マッピングは検出済みのはず");
}
let mut ports = req.ports().cloned().unwrap_or_default();
for &exposed in req.expose_ports() {
if ports.iter().any(|p| p.container_port() == exposed) {
continue;
}
ports.push(PortMapping::new(0, exposed));
}
crate::core::client::ContainerConfig {
image: req.descriptor(),
entrypoint: req.entrypoint().map(|e| vec![e.to_string()]),
cmd: req.cmd().map(|c| c.into_owned()).collect(),
env: req.env_vars().map(|(k, v)| format!("{k}={v}")).collect(),
ports,
mounts: req.mounts().cloned().collect(),
name: req.container_name().clone(),
labels: req.labels().clone(),
privileged: req.privileged(),
working_dir: req.working_dir().map(|d| d.to_string()),
user: req.user().map(|u| u.to_string()),
init: req.init(),
health_check: req.health_check().cloned(),
cap_add: req.cap_add().cloned().unwrap_or_default(),
cap_drop: req.cap_drop().cloned().unwrap_or_default(),
shm_size: req.shm_size(),
readonly_rootfs: req.readonly_rootfs(),
hostname: req.hostname().map(|h| h.to_string()),
open_stdin: req.open_stdin(),
network: req.network().clone(),
platform: req.platform().as_deref().map(|p| p.to_string()),
extra_hosts: req
.hosts()
.map(|(hostname, host)| format!("{hostname}:{host}"))
.collect(),
}
}
#[cfg(target_os = "macos")]
fn log_fd_required_error(cause: &crate::Error) -> crate::Error {
crate::Error::other(format!(
"log wait requires log file descriptors, but containerLogs failed: {cause}"
))
}
async fn rollback_remove(id: &str, remove_future: impl std::future::Future<Output = Result<()>>) {
if !matches!(
crate::core::env::Config.command(),
crate::core::env::Command::Remove
) {
return;
}
if let Err(rm_err) = remove_future.await {
tracing::warn!("failed to remove container {id} during rollback: {rm_err}");
}
}
#[cfg(target_os = "linux")]
fn linux_log_stream_required<I: Image>(
req: &ContainerRequest<I>,
ready_conditions: &[WaitFor],
) -> bool {
ready_conditions
.iter()
.any(|c| matches!(c, WaitFor::Log(_)))
|| !req.log_consumers.is_empty()
}
#[cfg(target_os = "linux")]
async fn start_linux_log_stream(
client: &crate::core::client::DockerClient,
id: &str,
log_consumers: Vec<Box<dyn crate::core::logs::consumer::LogConsumer + 'static>>,
log_required: bool,
) -> Result<(
crate::core::containers::async_container::ContainerLogSource,
Option<std::sync::Arc<Vec<Box<dyn crate::core::logs::consumer::LogConsumer + 'static>>>>,
)> {
use crate::core::client::docker_log_stream::spawn_log_consumer_task;
let has_consumers = !log_consumers.is_empty();
match client.spawn_log_session(id).await {
Ok(handle) => {
let stored = if has_consumers {
let consumers = std::sync::Arc::new(log_consumers);
spawn_log_consumer_task(
handle.clone(),
handle.stdout_stream(),
consumers.clone(),
crate::core::logs::LogFrame::StdOut,
);
spawn_log_consumer_task(
handle.clone(),
handle.stderr_stream(),
consumers.clone(),
crate::core::logs::LogFrame::StdErr,
);
Some(consumers)
} else {
None
};
Ok((ContainerLogSource::DockerStream(handle), stored))
}
Err(e) => {
if log_required {
rollback_remove(id, client.remove(id, true)).await;
return Err(e);
}
tracing::warn!("failed to start log stream; logs will be empty: {e}");
Ok((ContainerLogSource::None, None))
}
}
}
#[cfg(target_os = "linux")]
fn linux_unsupported_request_reason<I: Image>(req: &ContainerRequest<I>) -> Option<&'static str> {
if req.ssh() {
return Some("with_ssh() is not implemented on Linux");
}
if req.masked_paths().is_some() {
return Some("with_masked_paths() is not implemented on Linux");
}
if req.readonly_paths().is_some() {
return Some("with_readonly_paths() is not implemented on Linux");
}
None
}
#[cfg(target_os = "macos")]
async fn apply_extra_hosts<I: Image>(
container: &ContainerAsync<I>,
hosts: &[(String, ExtraHost)],
) -> Result<()> {
for (hostname, host) in hosts {
let ip = match host {
ExtraHost::Addr(ip) => ip.to_string(),
ExtraHost::HostGateway => {
let gw = container.gateway_ip_address().await.map_err(|e| {
crate::Error::other(format!(
"failed to resolve host gateway IP for with_host(..., HostGateway): {e}"
))
})?;
gw.to_string()
}
};
let cmd = ExecCommand::new([
"sh",
"-c",
r#"printf '%s\t%s\n' "$1" "$2" >> /etc/hosts"#,
"_",
&ip,
hostname,
])
.with_cmd_ready_condition(CmdWaitFor::exit_code(0));
container.exec(cmd).await?;
}
Ok(())
}
#[cfg(target_os = "macos")]
struct CopyDataTempFile {
path: std::path::PathBuf,
}
#[cfg(target_os = "macos")]
impl CopyDataTempFile {
fn as_path(&self) -> &std::path::Path {
&self.path
}
}
#[cfg(target_os = "macos")]
impl Drop for CopyDataTempFile {
fn drop(&mut self) {
if let Err(e) = std::fs::remove_file(&self.path) {
tracing::warn!(
"failed to remove copy_to temp file {}: {e}",
self.path.display()
);
}
}
}
#[cfg(target_os = "macos")]
async fn write_copy_data_temp(data: &[u8]) -> Result<CopyDataTempFile> {
use tokio::io::AsyncWriteExt;
let path =
std::env::temp_dir().join(format!("shiguredo_container_copy_{}.bin", unique_suffix()));
let mut file = tokio::fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(&path)
.await
.map_err(|e| crate::Error::other(format!("failed to create temp file for copy_to: {e}")))?;
let guard = CopyDataTempFile { path };
file.write_all(data)
.await
.map_err(|e| crate::Error::other(format!("failed to write temp file for copy_to: {e}")))?;
Ok(guard)
}
#[cfg(target_os = "macos")]
async fn copy_to_sources<I: Image>(
client: &crate::core::client::xpc_client::XpcClient,
id: &str,
req: &ContainerRequest<I>,
) -> Result<()> {
let sources: Vec<&CopyToContainer> = req.copy_to_sources().collect();
for src in sources {
match &src.source {
CopyDataSource::File(p) => {
client
.copy_in(id, p, &src.target.path, src.target.mode)
.await?;
chown_after_copy(
client,
id,
&src.target.path,
src.target.uid,
src.target.gid,
p.is_dir(),
)
.await;
}
CopyDataSource::Data(b) => {
let guard = write_copy_data_temp(b).await?;
client
.copy_in(id, guard.as_path(), &src.target.path, src.target.mode)
.await?;
chown_after_copy(
client,
id,
&src.target.path,
src.target.uid,
src.target.gid,
false,
)
.await;
}
}
}
Ok(())
}
#[cfg(target_os = "macos")]
async fn chown_after_copy(
client: &crate::core::client::xpc_client::XpcClient,
id: &str,
path: &str,
uid: u32,
gid: u32,
recursive: bool,
) {
if uid == 0 && gid == 0 {
return;
}
let mut cmd = vec!["chown".to_string()];
if recursive {
cmd.push("-R".to_string());
}
cmd.push(format!("{uid}:{gid}"));
cmd.push(path.to_string());
if let Err(e) = client.exec(id, &cmd, vec![]).await {
tracing::warn!("failed to chown {path} to {uid}:{gid}: {e}");
}
}
#[cfg(target_os = "linux")]
async fn copy_to_sources_linux<I: Image>(
client: &crate::core::client::docker_client::DockerClient,
id: &str,
req: &ContainerRequest<I>,
) -> Result<()> {
use std::path::{Path, PathBuf};
use crate::core::client::docker_client::DOCKER_RESPONSE_BODY_LIMIT;
use crate::core::client::docker_tar::UstarBuilder;
use crate::core::copy::CopyToContainerError;
fn size_limit_err(path: &Path) -> crate::Error {
crate::Error::other(CopyToContainerError::SizeLimitExceeded {
limit: DOCKER_RESPONSE_BODY_LIMIT,
name: path.display().to_string(),
})
}
fn check_file_size_limit(path: &Path, size: u64) -> Result<()> {
if size > DOCKER_RESPONSE_BODY_LIMIT as u64 {
return Err(size_limit_err(path));
}
Ok(())
}
fn read_file_limited(path: &Path) -> Result<Vec<u8>> {
use std::io::Read;
let file = std::fs::File::open(path)
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
let mut data = Vec::new();
let read = file
.take(DOCKER_RESPONSE_BODY_LIMIT as u64 + 1)
.read_to_end(&mut data)
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
if read > DOCKER_RESPONSE_BODY_LIMIT {
return Err(size_limit_err(path));
}
Ok(data)
}
async fn read_file_limited_async(path: &Path) -> Result<Vec<u8>> {
use tokio::io::AsyncReadExt;
let file = tokio::fs::File::open(path)
.await
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
let mut data = Vec::new();
let read = file
.take(DOCKER_RESPONSE_BODY_LIMIT as u64 + 1)
.read_to_end(&mut data)
.await
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
if read > DOCKER_RESPONSE_BODY_LIMIT {
return Err(size_limit_err(path));
}
Ok(data)
}
fn name_err(msg: impl Into<String>) -> crate::Error {
crate::Error::other(CopyToContainerError::PathNameError(msg.into()))
}
fn make_path_relative(path: &str) -> Result<String> {
let relative = path.trim_start_matches('/');
if relative.is_empty() {
return Err(name_err("copy_to target path must not be root only"));
}
Ok(relative.to_string())
}
fn append_ancestor_directories(
builder: &mut UstarBuilder,
relative: &str,
uid: u32,
gid: u32,
) -> Result<()> {
let mut acc = String::new();
let parts: Vec<&str> = relative.split('/').collect();
for part in parts.iter().take(parts.len().saturating_sub(1)) {
if part.is_empty() {
continue;
}
if !acc.is_empty() {
acc.push('/');
}
acc.push_str(part);
let dir_path = format!("{acc}/");
builder
.append_directory(&dir_path, 0o755, uid, gid)
.map_err(crate::Error::other)?;
}
Ok(())
}
fn append_host_directory(
builder: &mut UstarBuilder,
host_root: &Path,
tar_root: &str,
mode: u32,
uid: u32,
gid: u32,
) -> Result<()> {
let mut stack: Vec<PathBuf> = vec![host_root.to_path_buf()];
while let Some(dir) = stack.pop() {
let entries = std::fs::read_dir(&dir)
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
let mut children: Vec<_> = entries
.collect::<std::io::Result<Vec<_>>>()
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
children.sort_by_key(|e| e.file_name());
for entry in children.into_iter().rev() {
let path = entry.path();
let meta = std::fs::symlink_metadata(&path)
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
let ft = meta.file_type();
let rel = path
.strip_prefix(host_root)
.map_err(|_| name_err("copy_to walk path escaped source root"))?;
let rel_str = rel.to_str().ok_or_else(|| {
name_err("copy_to source path under directory is not valid UTF-8")
})?;
let tar_path = if rel_str.is_empty() {
tar_root.to_string()
} else {
format!("{tar_root}/{rel_str}")
};
if ft.is_symlink() {
return Err(name_err("copy_to source contains a symlink"));
}
if ft.is_dir() {
builder
.append_directory(&format!("{tar_path}/"), 0o755, uid, gid)
.map_err(crate::Error::other)?;
stack.push(path);
continue;
}
if !ft.is_file() {
return Err(name_err(
"copy_to source contains a non-regular file".to_string(),
));
}
check_file_size_limit(&path, meta.len())?;
let data = read_file_limited(&path)?;
builder
.append_file(&tar_path, &data, mode, uid, gid)
.map_err(crate::Error::other)?;
}
}
Ok(())
}
let sources: Vec<&CopyToContainer> = req.copy_to_sources().collect();
for src in sources {
let target_path = Path::new(&src.target.path);
if !target_path.is_absolute() {
return Err(name_err("copy_to target path must be absolute"));
}
if src.target.path.ends_with('/') {
return Err(name_err("copy_to target path must not end with a slash"));
}
if src.target.path.ends_with("/.") {
return Err(name_err("copy_to target path must not end with a '.'"));
}
if src.target.path.contains("//") {
return Err(name_err("copy_to target path must not contain '//'"));
}
if target_path.file_name().is_none() {
return Err(name_err("copy_to target path must have a file name"));
}
if target_path
.components()
.any(|c| matches!(c, std::path::Component::ParentDir))
{
return Err(name_err("copy_to target path must not contain '..'"));
}
let relative = make_path_relative(&src.target.path)?;
let mode = src.target.mode;
let uid = src.target.uid;
let gid = src.target.gid;
let mut builder = UstarBuilder::new();
match &src.source {
CopyDataSource::File(path) => {
let meta = tokio::fs::symlink_metadata(path)
.await
.map_err(|e| crate::Error::other(CopyToContainerError::IoError(e)))?;
let ft = meta.file_type();
if ft.is_symlink() {
return Err(name_err("copy_to source is a symlink"));
}
if ft.is_dir() {
append_ancestor_directories(&mut builder, &relative, uid, gid)?;
builder
.append_directory(&format!("{relative}/"), 0o755, uid, gid)
.map_err(crate::Error::other)?;
append_host_directory(&mut builder, path, &relative, mode, uid, gid)?;
} else if ft.is_file() {
append_ancestor_directories(&mut builder, &relative, uid, gid)?;
check_file_size_limit(path, meta.len())?;
let data = read_file_limited_async(path).await?;
builder
.append_file(&relative, &data, mode, uid, gid)
.map_err(crate::Error::other)?;
} else {
return Err(name_err(
"copy_to source is not a regular file or directory",
));
}
}
CopyDataSource::Data(data) => {
append_ancestor_directories(&mut builder, &relative, uid, gid)?;
builder
.append_file(&relative, data, mode, uid, gid)
.map_err(crate::Error::other)?;
}
}
let tar = builder.finish().map_err(crate::Error::other)?;
client.copy_to(id, "/", tar).await?;
}
Ok(())
}
#[cfg(all(test, target_os = "macos"))]
mod tests {
use super::*;
#[test]
fn unique_suffix_has_no_collision_under_parallel_generation() {
let handles: Vec<_> = (0..8)
.map(|_| std::thread::spawn(|| (0..1000).map(|_| unique_suffix()).collect::<Vec<_>>()))
.collect();
let mut all: Vec<String> = handles
.into_iter()
.flat_map(|h| h.join().expect("スレッドが panic しないこと"))
.collect();
let total = all.len();
all.sort();
all.dedup();
assert_eq!(all.len(), total, "unique_suffix should never collide");
}
#[cfg(target_os = "macos")]
#[test]
fn log_fd_required_error_includes_cause() {
let cause = crate::Error::other("simulated containerLogs failure");
let err = log_fd_required_error(&cause);
let msg = err.to_string();
assert!(
msg.contains("log wait requires log file descriptors, but containerLogs failed:"),
"明示メッセージが含まれること: {msg}"
);
assert!(
msg.contains("simulated containerLogs failure"),
"原因が含まれること: {msg}"
);
}
#[tokio::test]
async fn write_copy_data_temp_creates_file_with_mode_600() {
use std::os::unix::fs::PermissionsExt;
let guard = write_copy_data_temp(b"mode-check")
.await
.expect("一時ファイルの作成に失敗した");
let mode = std::fs::metadata(guard.as_path())
.expect("メタデータの取得に失敗した")
.permissions()
.mode();
assert_eq!(
mode & 0o777,
0o600,
"一時ファイルの mode が 0o600 であること"
);
}
}
#[cfg(all(test, target_os = "linux"))]
mod linux_tests {
use crate::core::ports::IntoContainerPort;
use crate::{ContainerRequest, GenericImage, ImageExt};
use super::{build_container_config, linux_unsupported_request_reason};
#[test]
fn masked_readonly_paths_are_unsupported_on_linux() {
let req: ContainerRequest<GenericImage> = GenericImage::new("alpine", "latest")
.with_cmd(["sleep", "1"])
.with_masked_paths(std::iter::empty::<String>());
assert_eq!(
linux_unsupported_request_reason(&req),
Some("with_masked_paths() is not implemented on Linux")
);
let req: ContainerRequest<GenericImage> = GenericImage::new("alpine", "latest")
.with_cmd(["sleep", "1"])
.with_readonly_paths(["/etc"]);
assert_eq!(
linux_unsupported_request_reason(&req),
Some("with_readonly_paths() is not implemented on Linux")
);
}
#[test]
fn expose_only_gets_host_port_zero() {
let req: ContainerRequest<GenericImage> = GenericImage::new("alpine", "latest")
.with_exposed_port(80.tcp())
.with_cmd(["sleep", "1"]);
let cfg = build_container_config(&req);
assert_eq!(cfg.ports.len(), 1, "expose のみなら 1 本であること");
assert_eq!(cfg.ports[0].container_port(), 80.tcp());
assert_eq!(
cfg.ports[0].host_port(),
0,
"create 前は Docker ランダム割当用に host_port が 0 であること"
);
}
#[test]
fn mapped_port_takes_precedence_over_same_expose() {
let req: ContainerRequest<GenericImage> = GenericImage::new("alpine", "latest")
.with_exposed_port(80.tcp())
.with_mapped_port(18080, 80.tcp())
.with_cmd(["sleep", "1"]);
let cfg = build_container_config(&req);
assert_eq!(cfg.ports.len(), 1, "重複は 1 本に収まること");
assert_eq!(cfg.ports[0].host_port(), 18080);
assert_eq!(cfg.ports[0].container_port(), 80.tcp());
}
#[test]
fn mapped_zero_takes_precedence_over_same_expose() {
let req: ContainerRequest<GenericImage> = GenericImage::new("alpine", "latest")
.with_exposed_port(80.tcp())
.with_mapped_port(0, 80.tcp())
.with_cmd(["sleep", "1"]);
let cfg = build_container_config(&req);
assert_eq!(cfg.ports.len(), 1, "重複は 1 本に収まること");
assert_eq!(cfg.ports[0].host_port(), 0);
assert_eq!(cfg.ports[0].container_port(), 80.tcp());
}
#[test]
fn different_protocol_same_number_keeps_both() {
let req: ContainerRequest<GenericImage> = GenericImage::new("alpine", "latest")
.with_exposed_port(80.udp())
.with_mapped_port(18080, 80.tcp())
.with_cmd(["sleep", "1"]);
let cfg = build_container_config(&req);
assert_eq!(cfg.ports.len(), 2, "異 proto は 2 本であること");
assert!(
cfg.ports
.iter()
.any(|p| p.container_port() == 80.tcp() && p.host_port() == 18080),
"tcp の明示マッピングが残ること"
);
assert!(
cfg.ports
.iter()
.any(|p| p.container_port() == 80.udp() && p.host_port() == 0),
"udp の expose が host_port 0 で追加されること"
);
}
#[test]
fn empty_expose_leaves_ports_empty() {
let req: ContainerRequest<GenericImage> =
GenericImage::new("alpine", "latest").with_cmd(["sleep", "1"]);
let cfg = build_container_config(&req);
assert!(cfg.ports.is_empty(), "ports が空であること");
}
#[test]
fn duplicate_expose_collapses_to_one_mapping() {
let req: ContainerRequest<GenericImage> = GenericImage::new("alpine", "latest")
.with_exposed_port(80.tcp())
.with_exposed_port(80.tcp())
.with_cmd(["sleep", "1"]);
let cfg = build_container_config(&req);
assert_eq!(cfg.ports.len(), 1, "二重 expose は 1 本に収まること");
assert_eq!(cfg.ports[0].host_port(), 0);
assert_eq!(cfg.ports[0].container_port(), 80.tcp());
}
}