use std::collections::BTreeMap;
use std::io::{Read, Write};
use std::os::unix::net::UnixStream;
use shiguredo_http11::{Request, Response};
use crate::core::client::http_decode::BodyLimit;
use crate::core::client::{ContainerConfig, ContainerSnapshot, HealthProbe, HealthStatus};
use crate::core::containers::request::PortMapping;
use crate::core::error::{ClientError, Result};
use crate::core::healthcheck::Healthcheck;
use crate::core::ports::{ContainerPort, Ports};
const DEFAULT_DOCKER_SOCKET: &str = "/var/run/docker.sock";
const EXEC_EXIT_CODE_BACKOFF_MILLIS: &[u64] = &[10, 50, 200, 500, 1000];
fn parse_exec_inspect_state(body: &[u8]) -> Result<(bool, Option<i64>)> {
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let value = parsed.value();
let running = value
.to_member("Running")
.ok()
.and_then(|m| m.optional())
.and_then(|v| bool::try_from(v).ok())
.unwrap_or(false);
let exit_code = value
.to_member("ExitCode")
.ok()
.and_then(|m| m.optional())
.and_then(|v| i64::try_from(v).ok());
Ok((running, exit_code))
}
fn resolve_exec_exit_code(states: &[(bool, Option<i64>)]) -> Option<i64> {
states
.iter()
.find(|(running, _)| !running)
.and_then(|(_, exit_code)| *exit_code)
}
const DOCKER_STREAM_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
pub(crate) const DOCKER_RESPONSE_BODY_LIMIT: usize = 64 * 1024 * 1024;
pub(crate) struct DockerExecResult {
pub(crate) exit_code: Option<i64>,
pub(crate) stdout: Vec<u8>,
pub(crate) stderr: Vec<u8>,
}
fn http11_err(e: impl std::fmt::Display) -> crate::core::error::Error {
ClientError::Other(e.to_string()).into()
}
#[derive(Debug, Clone)]
pub(crate) struct DockerClient {
socket_path: String,
}
impl DockerClient {
pub(crate) fn detect() -> Result<Self> {
Ok(Self {
socket_path: DEFAULT_DOCKER_SOCKET.to_string(),
})
}
pub(crate) async fn pull_image(&self, descriptor: &str, platform: Option<&str>) -> Result<()> {
let (image, tag) = split_pull_reference(descriptor);
let mut path = format!(
"/images/create?fromImage={}&tag={}",
percent_encode_component(image),
percent_encode_component(tag)
);
if let Some(platform) = platform {
path.push_str(&format!("&platform={}", percent_encode_component(platform)));
}
let auth_header = super::registry_auth::x_registry_auth(descriptor);
let extra_headers: Vec<(&'static str, String)> = auth_header
.map(|v| vec![("X-Registry-Auth", v)])
.unwrap_or_default();
let response = self
.request_with_extra_headers(
"POST",
&path,
None,
extra_headers,
BodyLimit::Error(DOCKER_RESPONSE_BODY_LIMIT),
)
.await
.map_err(|e| {
crate::core::error::Error::Client(ClientError::Other(format!(
"failed to receive pull progress for image {descriptor}: {e}"
)))
})?;
if response.status_code() >= 400 {
return Err(ClientError::Other(pull_error_message(
descriptor,
response.status_code(),
response.body_bytes().unwrap_or(&[]),
))
.into());
}
check_pull_stream_errors(response.body_bytes().unwrap_or(&[]))?;
Ok(())
}
pub(crate) async fn resolve_image_descriptor(
&self,
descriptor: &str,
platform: Option<&str>,
) -> Result<String> {
let path = format!("/images/{}/json", percent_encode_path_segment(descriptor));
let response = self.request("GET", &path, None).await?;
if response.status_code() == 200 {
Ok(descriptor.to_string())
} else if response.status_code() == 404 {
self.pull_image(descriptor, platform).await?;
let response = self.request("GET", &path, None).await?;
if response.status_code() == 200 {
Ok(descriptor.to_string())
} else {
Err(ClientError::ImageNotFound(descriptor.to_string()).into())
}
} else {
Err(ClientError::Other(format!(
"failed to resolve image {descriptor}: {}",
response.status_code()
))
.into())
}
}
pub(crate) async fn create_container(&self, config: ContainerConfig) -> Result<String> {
let mut params = Vec::new();
if let Some(name) = &config.name {
params.push(format!("name={}", percent_encode_component(name)));
}
if let Some(platform) = &config.platform {
params.push(format!("platform={}", percent_encode_component(platform)));
}
let query = if params.is_empty() {
String::new()
} else {
format!("?{}", params.join("&"))
};
let path = format!("/containers/create{query}");
let body = CreateContainerBody::from_config(config)?;
let body_json = body.to_json_string()?;
let response = self
.request("POST", &path, Some(body_json.into_bytes()))
.await?;
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to create container: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty response body".into()))?;
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let id: String = parsed
.value()
.to_member("Id")
.map_err(|e| ClientError::Json(e.to_string()))?
.required()
.map_err(|e| ClientError::Json(e.to_string()))?
.try_into()
.map_err(|e: nojson::JsonParseError| ClientError::Json(e.to_string()))?;
Ok(id)
}
pub(crate) async fn start_container(&self, id: &str) -> Result<()> {
let path = format!("/containers/{}/start", percent_encode_path_segment(id));
let response = self.request("POST", &path, None).await?;
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to start container: {}",
response.status_code()
))
.into());
}
Ok(())
}
pub(crate) async fn stop(&self, id: &str, timeout_seconds: Option<i32>) -> Result<()> {
let t = match timeout_seconds {
Some(t) if t >= 0 => t,
_ => 30,
};
let path = format!("/containers/{}/stop?t={t}", percent_encode_path_segment(id));
let response = self.request("POST", &path, None).await?;
if response.status_code() == 404 {
return Ok(());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to stop container: {}",
response.status_code()
))
.into());
}
Ok(())
}
pub(crate) async fn pause(&self, id: &str) -> Result<()> {
let path = format!("/containers/{}/pause", percent_encode_path_segment(id));
let response = self.request("POST", &path, None).await?;
if response.status_code() == 304 || response.status_code() == 404 {
return Ok(());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to pause container: {}",
response.status_code()
))
.into());
}
Ok(())
}
pub(crate) async fn unpause(&self, id: &str) -> Result<()> {
let path = format!("/containers/{}/unpause", percent_encode_path_segment(id));
let response = self.request("POST", &path, None).await?;
if response.status_code() == 304 || response.status_code() == 404 {
return Ok(());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to unpause container: {}",
response.status_code()
))
.into());
}
Ok(())
}
pub(crate) async fn remove(&self, id: &str, force: bool) -> Result<()> {
let path = format!(
"/containers/{}?force={force}",
percent_encode_path_segment(id)
);
let response = self.request("DELETE", &path, None).await?;
if response.status_code() == 404 {
return Ok(());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to remove container: {}",
response.status_code()
))
.into());
}
Ok(())
}
pub(crate) fn remove_blocking(&self, id: &str, force: bool) -> Result<()> {
let path = format!(
"/containers/{}?force={force}",
percent_encode_path_segment(id)
);
let request_bytes = encode_docker_api_request("DELETE", &path, None)?;
let mut stream = UnixStream::connect(&self.socket_path)?;
stream.set_read_timeout(Some(DOCKER_STREAM_TIMEOUT))?;
stream.set_write_timeout(Some(DOCKER_STREAM_TIMEOUT))?;
stream.write_all(&request_bytes)?;
let response = read_http11_response(&mut stream, "DELETE", BodyLimit::Unlimited)?;
if response.status_code() == 404 {
return Ok(());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to remove container: {}",
response.status_code()
))
.into());
}
Ok(())
}
pub(crate) fn wait_blocking(&self, id: &str) -> Result<i64> {
let path = format!(
"/containers/{}/wait?condition=not-running",
percent_encode_path_segment(id)
);
let request_bytes = encode_docker_api_request("POST", &path, None)?;
let mut stream = UnixStream::connect(&self.socket_path)?;
stream.write_all(&request_bytes)?;
let response = read_http11_response(&mut stream, "POST", BodyLimit::Unlimited)?;
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to wait for container: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty wait response body".into()))?;
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let status_code = parsed
.value()
.to_member("StatusCode")
.map_err(|e| ClientError::Json(e.to_string()))?
.required()
.map_err(|e| ClientError::Json(e.to_string()))?;
let code: i64 = status_code
.try_into()
.map_err(|e: nojson::JsonParseError| ClientError::Json(e.to_string()))?;
Ok(code)
}
pub(crate) async fn exec(
&self,
id: &str,
cmd: &[String],
env: Vec<String>,
) -> Result<DockerExecResult> {
let exec_path = format!("/containers/{}/exec", percent_encode_path_segment(id));
let exec_config = ExecConfig {
cmd: cmd.to_vec(),
attach_stdout: true,
attach_stderr: true,
env,
};
let exec_json = exec_config.to_json_string()?;
let response = self
.request("POST", &exec_path, Some(exec_json.into_bytes()))
.await?;
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to create exec: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty exec response body".into()))?;
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let exec_id: String = parsed
.value()
.to_member("Id")
.map_err(|e| ClientError::Json(e.to_string()))?
.required()
.map_err(|e| ClientError::Json(e.to_string()))?
.try_into()
.map_err(|e: nojson::JsonParseError| ClientError::Json(e.to_string()))?;
let start_path = format!("/exec/{}/start", percent_encode_path_segment(&exec_id));
let start_config = ExecStartConfig {
detach: false,
tty: false,
};
let start_json = start_config.to_json_string()?;
let response = self
.request_with_body_limit(
"POST",
&start_path,
Some(start_json.into_bytes()),
BodyLimit::Error(DOCKER_RESPONSE_BODY_LIMIT),
)
.await?;
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to start exec: {}",
response.status_code()
))
.into());
}
let stream_body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty exec start response body".into()))?;
let (stdout, stderr) = demux_exec_stream(stream_body);
let inspect_path = format!("/exec/{}/json", percent_encode_path_segment(&exec_id));
let mut states: Vec<(bool, Option<i64>)> = Vec::new();
#[expect(clippy::needless_range_loop)]
for attempt in 0..=EXEC_EXIT_CODE_BACKOFF_MILLIS.len() {
let response = self.request("GET", &inspect_path, None).await?;
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to inspect exec: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty exec inspect body".into()))?;
let state = parse_exec_inspect_state(body)?;
let (running, _) = state;
states.push(state);
if !running {
break;
}
if attempt < EXEC_EXIT_CODE_BACKOFF_MILLIS.len() {
tokio::time::sleep(std::time::Duration::from_millis(
EXEC_EXIT_CODE_BACKOFF_MILLIS[attempt],
))
.await;
}
}
let exit_code = resolve_exec_exit_code(&states);
if exit_code.is_none() {
tracing::warn!("exec {exec_id} still running after stream EOF, exit code unavailable");
}
Ok(DockerExecResult {
exit_code,
stdout,
stderr,
})
}
pub(crate) async fn ports(&self, id: &str) -> Result<Ports> {
let snapshot = self.container_state(id).await?;
Ok(snapshot.ports)
}
pub(crate) async fn bridge_ip_address(&self, id: &str) -> Result<std::net::IpAddr> {
let path = format!("/containers/{}/json", percent_encode_path_segment(id));
let response = self.request("GET", &path, None).await?;
if response.status_code() == 404 {
return Err(ClientError::ContainerNotFound(id.to_string()).into());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to inspect container: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty inspect body".into()))?;
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let ip_str = parsed
.value()
.to_member("NetworkSettings")
.ok()
.and_then(|m| m.optional())
.and_then(|ns| ns.to_member("Networks").ok())
.and_then(|m| m.optional())
.and_then(|networks| networks.to_object().ok())
.and_then(|mut obj| obj.next())
.and_then(|(_, network)| {
network
.to_member("IPAddress")
.ok()
.and_then(|m| m.optional())
.and_then(|v| String::try_from(v).ok())
})
.filter(|s| !s.is_empty())
.ok_or_else(|| {
ClientError::Other("no network IP address found for container".into())
})?;
ip_str
.parse::<std::net::IpAddr>()
.map_err(|e| ClientError::Other(format!("invalid IP address '{ip_str}': {e}")).into())
}
pub(crate) async fn container_state(&self, id: &str) -> Result<ContainerSnapshot> {
let path = format!("/containers/{}/json", percent_encode_path_segment(id));
let response = self.request("GET", &path, None).await?;
if response.status_code() == 404 {
return Err(ClientError::ContainerNotFound(id.to_string()).into());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to inspect container: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty inspect body".into()))?;
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let state = parsed
.value()
.to_member("State")
.map_err(|e| ClientError::Json(e.to_string()))?
.required()
.map_err(|e| ClientError::Json(e.to_string()))?;
let running = state
.to_member("Running")
.ok()
.and_then(|m| m.optional())
.and_then(|v| bool::try_from(v).ok())
.unwrap_or(false);
let ports = parse_ports(&parsed);
Ok(ContainerSnapshot { running, ports })
}
pub(crate) async fn container_env(&self, id: &str) -> Result<Vec<String>> {
let path = format!("/containers/{}/json", percent_encode_path_segment(id));
let response = self.request("GET", &path, None).await?;
if response.status_code() == 404 {
return Err(ClientError::ContainerNotFound(id.to_string()).into());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to inspect container: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty inspect body".into()))?;
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let mut env = Vec::new();
if let Ok(config) = parsed.value().to_member("Config")
&& let Some(config) = config.optional()
&& let Ok(env_member) = config.to_member("Env")
&& let Some(env_arr) = env_member.optional()
&& let Ok(arr) = env_arr.to_array()
{
for item in arr {
if let Ok(s) = String::try_from(item) {
env.push(s);
}
}
}
Ok(env)
}
pub(crate) async fn container_health(&self, id: &str) -> Result<HealthProbe> {
let path = format!("/containers/{}/json", percent_encode_path_segment(id));
let response = self.request("GET", &path, None).await?;
if response.status_code() == 404 {
return Err(ClientError::ContainerNotFound(id.to_string()).into());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to inspect container: {}",
response.status_code()
))
.into());
}
let body = response
.body_bytes()
.ok_or_else(|| ClientError::Other("empty inspect body".into()))?;
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
let parsed = nojson::RawJson::parse(text).map_err(|e| ClientError::Json(e.to_string()))?;
let state = parsed
.value()
.to_member("State")
.map_err(|e| ClientError::Json(e.to_string()))?
.required()
.map_err(|e| ClientError::Json(e.to_string()))?;
let running = state
.to_member("Running")
.ok()
.and_then(|m| m.optional())
.and_then(|v| bool::try_from(v).ok())
.unwrap_or(false);
let health = state
.to_member("Health")
.ok()
.and_then(|m| m.optional())
.and_then(|health| {
health
.to_member("Status")
.ok()
.and_then(|m| m.optional())
.and_then(|v| String::try_from(v).ok())
})
.and_then(|status| match status.as_str() {
"starting" => Some(HealthStatus::Starting),
"healthy" => Some(HealthStatus::Healthy),
"unhealthy" => Some(HealthStatus::Unhealthy),
"none" | "" => None,
_ => None,
});
Ok(HealthProbe { running, health })
}
pub(crate) async fn spawn_log_session(
&self,
id: &str,
) -> Result<std::sync::Arc<crate::core::client::docker_log_stream::DockerLogsHandle>> {
crate::core::client::docker_log_stream::spawn_log_session(
self.socket_path.clone(),
id.to_string(),
)
.await
}
async fn request(&self, method: &str, path: &str, body: Option<Vec<u8>>) -> Result<Response> {
self.request_with_body_limit(method, path, body, BodyLimit::Unlimited)
.await
}
async fn request_with_body_limit(
&self,
method: &str,
path: &str,
body: Option<Vec<u8>>,
body_limit: BodyLimit,
) -> Result<Response> {
let socket_path = self.socket_path.clone();
let method = method.to_string();
let path = path.to_string();
tokio::task::spawn_blocking(move || -> Result<Response> {
let request_bytes = encode_docker_api_request(
&method,
&path,
body.as_deref().map(|b| (b, "application/json")),
)?;
let mut stream = UnixStream::connect(&socket_path)?;
stream.write_all(&request_bytes)?;
read_http11_response(&mut stream, &method, body_limit)
})
.await
.map_err(|e| ClientError::Other(format!("spawn_blocking failed: {e}")))?
}
async fn request_with_extra_headers(
&self,
method: &str,
path: &str,
body: Option<Vec<u8>>,
extra_headers: Vec<(&'static str, String)>,
body_limit: BodyLimit,
) -> Result<Response> {
let socket_path = self.socket_path.clone();
let method = method.to_string();
let path = path.to_string();
tokio::task::spawn_blocking(move || -> Result<Response> {
let request_bytes = encode_docker_api_request_with_headers(
&method,
&path,
body.as_deref().map(|b| (b, "application/json")),
&extra_headers,
)?;
let mut stream = UnixStream::connect(&socket_path)?;
stream.write_all(&request_bytes)?;
read_http11_response(&mut stream, &method, body_limit)
})
.await
.map_err(|e| ClientError::Other(format!("spawn_blocking failed: {e}")))?
}
async fn request_with_content_type(
&self,
method: &str,
path: &str,
body: Vec<u8>,
content_type: &'static str,
) -> Result<Response> {
let socket_path = self.socket_path.clone();
let method = method.to_string();
let path = path.to_string();
tokio::task::spawn_blocking(move || -> Result<Response> {
let request_bytes =
encode_docker_api_request(&method, &path, Some((body.as_slice(), content_type)))?;
let mut stream = UnixStream::connect(&socket_path)?;
stream.write_all(&request_bytes)?;
read_http11_response(&mut stream, &method, BodyLimit::Unlimited)
})
.await
.map_err(|e| ClientError::Other(format!("spawn_blocking failed: {e}")))?
}
pub(crate) async fn copy_from(&self, id: &str, path: &str) -> Result<Vec<u8>> {
let api_path = format!(
"/containers/{}/archive?path={}",
percent_encode_path_segment(id),
percent_encode_component(path)
);
let response = self
.request_with_body_limit(
"GET",
&api_path,
None,
BodyLimit::Error(DOCKER_RESPONSE_BODY_LIMIT),
)
.await?;
if response.status_code() == 404 {
return Err(
classify_archive_404(id, path, response.body_bytes().unwrap_or(&[])).into(),
);
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to copy from container: {}",
response.status_code()
))
.into());
}
Ok(response.body_bytes().unwrap_or(&[]).to_vec())
}
pub(crate) async fn copy_to(&self, id: &str, dir: &str, tar: Vec<u8>) -> Result<()> {
let api_path = format!(
"/containers/{}/archive?path={}©UIDGID=true",
percent_encode_path_segment(id),
percent_encode_component(dir)
);
let response = self
.request_with_content_type("PUT", &api_path, tar, "application/x-tar")
.await?;
if response.status_code() == 404 {
return Err(ClientError::ContainerNotFound(id.to_string()).into());
}
if response.status_code() >= 400 {
return Err(ClientError::Other(format!(
"failed to copy to container: {}",
response.status_code()
))
.into());
}
Ok(())
}
}
fn classify_archive_404(id: &str, path: &str, body: &[u8]) -> ClientError {
let message = parse_daemon_error_message(body);
match message {
Some(msg) if msg.starts_with("Could not find the file ") => {
ClientError::ContainerPathNotFound(path.to_string())
}
_ => ClientError::ContainerNotFound(id.to_string()),
}
}
fn parse_daemon_error_message(body: &[u8]) -> Option<String> {
let text = std::str::from_utf8(body).ok()?;
let parsed = nojson::RawJson::parse(text).ok()?;
parsed
.value()
.to_member("message")
.ok()
.and_then(|m| m.required().ok())
.and_then(|v| TryInto::<String>::try_into(v).ok())
}
fn pull_error_message(descriptor: &str, status: u16, body: &[u8]) -> String {
match parse_daemon_error_message(body) {
Some(message) if !message.is_empty() => {
format!("failed to pull image {descriptor}: {status}: {message}")
}
_ => format!("failed to pull image {descriptor}: {status}"),
}
}
pub(crate) fn encode_docker_api_request(
method: &str,
path: &str,
body: Option<(&[u8], &'static str)>,
) -> Result<Vec<u8>> {
let method_obj = shiguredo_http11::Method::new(method).map_err(http11_err)?;
let mut request = Request::new(method_obj, path)
.map_err(http11_err)?
.header("Host", "localhost")
.map_err(http11_err)?
.header("Connection", "close")
.map_err(http11_err)?;
if let Some((body, content_type)) = body {
request = request
.header("Content-Type", content_type)
.map_err(http11_err)?;
request = request
.header("Content-Length", &body.len().to_string())
.map_err(http11_err)?;
request = request.body(body.to_vec());
}
request.encode().map_err(http11_err)
}
fn encode_docker_api_request_with_headers(
method: &str,
path: &str,
body: Option<(&[u8], &'static str)>,
extra_headers: &[(&'static str, String)],
) -> Result<Vec<u8>> {
let method_obj = shiguredo_http11::Method::new(method).map_err(http11_err)?;
let mut request = Request::new(method_obj, path)
.map_err(http11_err)?
.header("Host", "localhost")
.map_err(http11_err)?
.header("Connection", "close")
.map_err(http11_err)?;
for (name, value) in extra_headers {
request = request.header(*name, value.as_str()).map_err(http11_err)?;
}
if let Some((body, content_type)) = body {
request = request
.header("Content-Type", content_type)
.map_err(http11_err)?;
request = request
.header("Content-Length", &body.len().to_string())
.map_err(http11_err)?;
request = request.body(body.to_vec());
}
request.encode().map_err(http11_err)
}
fn read_http11_response(
stream: &mut impl Read,
method: &str,
body_limit: BodyLimit,
) -> Result<Response> {
use crate::core::client::http_decode::ResponseAccumulator;
let mut acc = ResponseAccumulator::new(method, body_limit);
loop {
let want = acc.read_buf_size();
if want == 0 {
return Err(ClientError::Other("decoder buffer full".into()).into());
}
let buf = acc.mut_buf(want).map_err(ClientError::Other)?;
let n = stream.read(buf)?;
if acc.feed(n).map_err(ClientError::Other)? {
break;
}
}
let decoded = acc.finish().map_err(ClientError::Other)?;
let mut response = Response::with_version(
decoded.head.version(),
decoded.head.status_code(),
decoded.head.reason_phrase(),
)
.map_err(http11_err)?;
for (name, value) in decoded.head.headers() {
let name = shiguredo_http11::HeaderName::new(name.as_str()).map_err(http11_err)?;
response = response.header(name, value).map_err(http11_err)?;
}
response = response.body(decoded.body);
Ok(response)
}
struct CreateContainerBody {
image: String,
entrypoint: Option<Vec<String>>,
cmd: Vec<String>,
env: Vec<String>,
labels: BTreeMap<String, String>,
working_dir: Option<String>,
user: Option<String>,
host_config: HostConfig,
exposed_ports: Vec<String>,
healthcheck: Option<Healthcheck>,
hostname: Option<String>,
open_stdin: Option<bool>,
network: Option<String>,
}
impl CreateContainerBody {
fn from_config(config: ContainerConfig) -> Result<Self> {
let port_bindings = build_port_bindings(&config.ports);
let exposed_ports = build_exposed_ports(&config.ports);
let mut binds = Vec::new();
let mut mounts = Vec::new();
for m in &config.mounts {
match m.mount_type() {
crate::core::mounts::MountType::Bind => {
let source = m.source().ok_or_else(|| {
ClientError::Configuration("bind mount source is required".into())
})?;
let target = m.target().ok_or_else(|| {
ClientError::Configuration("bind mount target is required".into())
})?;
binds.push(format!("{source}:{target}:{}", m.access_mode()));
}
crate::core::mounts::MountType::Volume => {
let source = m.source().ok_or_else(|| {
ClientError::Configuration("volume mount source is required".into())
})?;
let target = m.target().ok_or_else(|| {
ClientError::Configuration("volume mount target is required".into())
})?;
let read_only = m.access_mode() == crate::core::mounts::AccessMode::ReadOnly;
mounts.push(format!(
"{{\"Type\":\"volume\",\"Source\":{},\"Target\":{},\"ReadOnly\":{}}}",
escape_json(source),
escape_json(target),
read_only
));
}
crate::core::mounts::MountType::Tmpfs => {
let target = m.target().ok_or_else(|| {
ClientError::Configuration("tmpfs mount target is required".into())
})?;
let read_only = m.access_mode() == crate::core::mounts::AccessMode::ReadOnly;
let mut entry = format!(
"{{\"Type\":\"tmpfs\",\"Target\":{},\"ReadOnly\":{}",
escape_json(target),
read_only
);
if let Some(opts) = m.tmpfs_options() {
let mut tmpfs_opts = String::new();
if let Some(size) = opts.size_bytes() {
tmpfs_opts.push_str(&format!("\"SizeBytes\":{size}"));
}
if let Some(mode) = opts.mode() {
if !tmpfs_opts.is_empty() {
tmpfs_opts.push(',');
}
tmpfs_opts.push_str(&format!("\"Mode\":{mode}"));
}
if !tmpfs_opts.is_empty() {
entry.push_str(&format!(",\"TmpfsOptions\":{{{tmpfs_opts}}}"));
}
}
entry.push('}');
mounts.push(entry);
}
}
}
Ok(Self {
image: config.image,
entrypoint: config.entrypoint,
cmd: config.cmd,
env: config.env,
labels: config.labels,
working_dir: config.working_dir,
user: config.user,
host_config: HostConfig {
port_bindings,
binds,
privileged: config.privileged,
init: config.init,
cap_add: config.cap_add,
cap_drop: config.cap_drop,
shm_size: config.shm_size,
readonly_rootfs: config.readonly_rootfs,
extra_hosts: config.extra_hosts,
mounts,
},
exposed_ports,
healthcheck: config.health_check,
hostname: config.hostname,
open_stdin: config.open_stdin,
network: config.network,
})
}
fn to_json_string(&self) -> Result<String> {
let mut json = String::new();
json.push_str("{\"Image\":");
json.push_str(&escape_json(&self.image));
if let Some(entrypoint) = &self.entrypoint
&& !entrypoint.is_empty()
{
json.push_str(",\"Entrypoint\":");
json.push_str(&json_array(entrypoint));
}
if !self.cmd.is_empty() {
json.push_str(",\"Cmd\":");
json.push_str(&json_array(&self.cmd));
}
if !self.env.is_empty() {
json.push_str(",\"Env\":");
json.push_str(&json_array(&self.env));
}
if !self.labels.is_empty() {
json.push_str(",\"Labels\":");
json.push_str(&json_object(&self.labels));
}
if let Some(wd) = &self.working_dir {
json.push_str(",\"WorkingDir\":");
json.push_str(&escape_json(wd));
}
if let Some(user) = &self.user {
json.push_str(",\"User\":");
json.push_str(&escape_json(user));
}
if !self.exposed_ports.is_empty() {
json.push_str(",\"ExposedPorts\":");
json.push('{');
for (i, key) in self.exposed_ports.iter().enumerate() {
if i > 0 {
json.push(',');
}
json.push_str(&escape_json(key));
json.push_str(":{}");
}
json.push('}');
}
if let Some(hc_json) = self.healthcheck.as_ref().and_then(|hc| hc.to_docker_json()) {
json.push_str(",\"Healthcheck\":");
json.push_str(&hc_json);
}
if let Some(hostname) = &self.hostname {
json.push_str(",\"Hostname\":");
json.push_str(&escape_json(hostname));
}
if self.open_stdin == Some(true) {
json.push_str(",\"OpenStdin\":true");
}
if let Some(network) = &self.network {
json.push_str(",\"NetworkingConfig\":{\"EndpointsConfig\":{");
json.push_str(&escape_json(network));
json.push_str(":{}");
json.push_str("}}");
}
json.push_str(",\"HostConfig\":");
json.push_str(&self.host_config.to_json_string()?);
json.push('}');
Ok(json)
}
}
struct HostConfig {
port_bindings: BTreeMap<String, Vec<PortBinding>>,
binds: Vec<String>,
privileged: bool,
init: bool,
cap_add: Vec<String>,
cap_drop: Vec<String>,
shm_size: Option<u64>,
readonly_rootfs: bool,
extra_hosts: Vec<String>,
mounts: Vec<String>,
}
impl HostConfig {
fn to_json_string(&self) -> Result<String> {
let mut json = String::new();
json.push_str("{\"Privileged\":");
json.push_str(if self.privileged { "true" } else { "false" });
json.push_str(",\"Init\":");
json.push_str(if self.init { "true" } else { "false" });
json.push_str(",\"ReadonlyRootfs\":");
json.push_str(if self.readonly_rootfs {
"true"
} else {
"false"
});
if !self.binds.is_empty() {
json.push_str(",\"Binds\":");
json.push_str(&json_array(&self.binds));
}
if !self.port_bindings.is_empty() {
json.push_str(",\"PortBindings\":");
json.push_str(&json_port_bindings(&self.port_bindings));
}
if !self.cap_add.is_empty() {
json.push_str(",\"CapAdd\":");
json.push_str(&json_array(&self.cap_add));
}
if !self.cap_drop.is_empty() {
json.push_str(",\"CapDrop\":");
json.push_str(&json_array(&self.cap_drop));
}
if let Some(shm_size) = self.shm_size {
json.push_str(",\"ShmSize\":");
json.push_str(&shm_size.to_string());
}
if !self.extra_hosts.is_empty() {
json.push_str(",\"ExtraHosts\":");
json.push_str(&json_array(&self.extra_hosts));
}
if !self.mounts.is_empty() {
json.push_str(",\"Mounts\":[");
json.push_str(&self.mounts.join(","));
json.push(']');
}
json.push('}');
Ok(json)
}
}
struct PortBinding {
host_ip: String,
host_port: String,
}
struct ExecConfig {
cmd: Vec<String>,
attach_stdout: bool,
attach_stderr: bool,
env: Vec<String>,
}
impl ExecConfig {
fn to_json_string(&self) -> Result<String> {
let mut json = String::new();
json.push_str("{\"AttachStdout\":");
json.push_str(if self.attach_stdout { "true" } else { "false" });
json.push_str(",\"AttachStderr\":");
json.push_str(if self.attach_stderr { "true" } else { "false" });
if !self.cmd.is_empty() {
json.push_str(",\"Cmd\":");
json.push_str(&json_array(&self.cmd));
}
if !self.env.is_empty() {
json.push_str(",\"Env\":");
json.push_str(&json_array(&self.env));
}
json.push('}');
Ok(json)
}
}
struct ExecStartConfig {
detach: bool,
tty: bool,
}
impl ExecStartConfig {
fn to_json_string(&self) -> Result<String> {
Ok(format!(
"{{\"Detach\":{},\"Tty\":{}}}",
if self.detach { "true" } else { "false" },
if self.tty { "true" } else { "false" }
))
}
}
pub(crate) fn percent_encode_path_segment(s: &str) -> String {
percent_encode(s)
}
fn demux_exec_stream(data: &[u8]) -> (Vec<u8>, Vec<u8>) {
const FRAME_HEADER_LEN: usize = 8;
let mut stdout = Vec::new();
let mut stderr = Vec::new();
let mut pos = 0;
while pos + FRAME_HEADER_LEN <= data.len() {
let stream_type = data[pos];
let payload_len =
u32::from_be_bytes([data[pos + 4], data[pos + 5], data[pos + 6], data[pos + 7]])
as usize;
pos += FRAME_HEADER_LEN;
let end = (pos + payload_len).min(data.len());
let chunk = &data[pos..end];
match stream_type {
1 => stdout.extend_from_slice(chunk),
2 => stderr.extend_from_slice(chunk),
_ => {}
}
pos = end;
}
(stdout, stderr)
}
fn percent_encode_component(s: &str) -> String {
percent_encode(s)
}
fn percent_encode(s: &str) -> String {
let mut out = String::new();
for b in s.bytes() {
match b {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
out.push(b as char);
}
_ => {
out.push('%');
const HEX: &[u8; 16] = b"0123456789ABCDEF";
out.push(HEX[(b >> 4) as usize] as char);
out.push(HEX[(b & 0xf) as usize] as char);
}
}
}
out
}
fn split_pull_reference(descriptor: &str) -> (&str, &str) {
if let Some((left, digest)) = descriptor.split_once('@') {
let tag = if digest.is_empty() { "latest" } else { digest };
let (from_image, _) = split_name_and_tag(left);
return (from_image, tag);
}
split_name_and_tag(descriptor)
}
fn split_name_and_tag(descriptor: &str) -> (&str, &str) {
let name_tag = match descriptor.rfind('/') {
Some(i) => &descriptor[i + 1..],
None => descriptor,
};
match name_tag.split_once(':') {
None => (descriptor, "latest"),
Some((name, "")) => {
let from_image = &descriptor[..descriptor.len() - name_tag.len() + name.len()];
(from_image, "latest")
}
Some((name, tag)) => {
let from_image = &descriptor[..descriptor.len() - name_tag.len() + name.len()];
(from_image, tag)
}
}
}
fn check_pull_stream_errors(body: &[u8]) -> std::result::Result<(), ClientError> {
let text = std::str::from_utf8(body).map_err(|e| ClientError::Json(e.to_string()))?;
for line in text.lines() {
let line = line.trim_end_matches('\r').trim();
if line.is_empty() {
continue;
}
let parsed = nojson::RawJson::parse(line).map_err(|e| ClientError::Json(e.to_string()))?;
if !parsed.value().kind().is_object() {
return Err(ClientError::Json(format!(
"expected object in pull stream, got {line}"
)));
}
let error = parsed
.value()
.to_member("error")
.ok()
.and_then(|m| m.optional());
let detail = parsed
.value()
.to_member("errorDetail")
.ok()
.and_then(|m| m.optional());
let error_msg = error
.and_then(|v| String::try_from(v).ok())
.filter(|s| !s.is_empty());
let detail_msg = detail
.and_then(|d| d.to_member("message").ok())
.and_then(|m| m.optional())
.and_then(|v| String::try_from(v).ok())
.filter(|s| !s.is_empty());
if let Some(msg) = error_msg.or(detail_msg) {
return Err(ClientError::Other(msg));
}
}
Ok(())
}
fn build_port_bindings(ports: &[PortMapping]) -> BTreeMap<String, Vec<PortBinding>> {
let mut map = BTreeMap::new();
for p in ports {
let key = format!(
"{}/{}",
p.container_port.as_u16(),
p.container_port.as_str()
);
map.insert(
key,
vec![PortBinding {
host_ip: "0.0.0.0".to_string(),
host_port: p.host_port.to_string(),
}],
);
}
map
}
fn build_exposed_ports(ports: &[PortMapping]) -> Vec<String> {
ports
.iter()
.map(|p| {
format!(
"{}/{}",
p.container_port.as_u16(),
p.container_port.as_str()
)
})
.collect()
}
fn parse_ports(parsed: &nojson::RawJson<'_>) -> Ports {
let mut ports = Ports::default();
let network_settings = match parsed.value().to_member("NetworkSettings") {
Ok(m) => match m.optional() {
Some(v) => v,
None => return ports,
},
Err(_) => return ports,
};
let ports_json = match network_settings.to_member("Ports") {
Ok(m) => match m.optional() {
Some(v) => v,
None => return ports,
},
Err(_) => return ports,
};
let Ok(entries) = ports_json.to_object() else {
return ports;
};
for (key, value) in entries {
let key_str: String = match key.try_into() {
Ok(s) => s,
Err(_) => continue,
};
let parts: Vec<&str> = key_str.split('/').collect();
if parts.len() != 2 {
continue;
}
let Ok(container_port) = parts[0].parse::<u16>() else {
continue;
};
let container_port = match parts[1] {
"tcp" => ContainerPort::Tcp(container_port),
"udp" => ContainerPort::Udp(container_port),
"sctp" => ContainerPort::Sctp(container_port),
_ => continue,
};
if let Ok(bindings) = value.to_array() {
for binding in bindings {
if let Some(host_port) = binding
.to_member("HostPort")
.ok()
.and_then(|m| m.optional())
.and_then(|v| String::try_from(v).ok())
.and_then(|s| s.parse().ok())
{
ports.add_mapping(container_port, host_port);
}
}
}
}
ports
}
pub(crate) fn json_array(items: &[String]) -> String {
let mut json = String::from("[");
for (i, item) in items.iter().enumerate() {
if i > 0 {
json.push(',');
}
json.push_str(&escape_json(item));
}
json.push(']');
json
}
fn json_object(map: &BTreeMap<String, String>) -> String {
let mut json = String::from("{");
for (i, (k, v)) in map.iter().enumerate() {
if i > 0 {
json.push(',');
}
json.push_str(&escape_json(k));
json.push(':');
json.push_str(&escape_json(v));
}
json.push('}');
json
}
fn json_port_bindings(map: &BTreeMap<String, Vec<PortBinding>>) -> String {
let mut json = String::from("{");
for (i, (k, bindings)) in map.iter().enumerate() {
if i > 0 {
json.push(',');
}
json.push_str(&escape_json(k));
json.push_str(":[");
for (j, b) in bindings.iter().enumerate() {
if j > 0 {
json.push(',');
}
json.push_str("{\"HostIp\":");
json.push_str(&escape_json(&b.host_ip));
json.push_str(",\"HostPort\":");
json.push_str(&escape_json(&b.host_port));
json.push('}');
}
json.push(']');
}
json.push('}');
json
}
pub(crate) fn escape_json(s: &str) -> String {
let mut escaped = String::with_capacity(s.len() + 2);
escaped.push('"');
for c in s.chars() {
match c {
'"' => escaped.push_str("\\\""),
'\\' => escaped.push_str("\\\\"),
'\u{0008}' => escaped.push_str("\\b"),
'\u{000c}' => escaped.push_str("\\f"),
'\n' => escaped.push_str("\\n"),
'\r' => escaped.push_str("\\r"),
'\t' => escaped.push_str("\\t"),
c if (c as u32) < 0x20 => escaped.push_str(&format!("\\u{:04x}", c as u32)),
c => escaped.push(c),
}
}
escaped.push('"');
escaped
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::ports::IntoContainerPort;
#[test]
fn escape_json_escapes_special_characters() {
let escaped = escape_json("a\"b\\c\nd\te\x01");
assert_eq!(
escaped, "\"a\\\"b\\\\c\\nd\\te\\u0001\"",
"引用符・バックスラッシュ・改行・タブ・制御文字が正しくエスケープされること"
);
}
#[test]
fn escape_json_roundtrips_through_nojson_for_samples() {
for s in ["", "plain", "a\"b", "line\n", "タブ\t", "\\"] {
let escaped = escape_json(s);
let parsed = nojson::RawJson::parse(&escaped)
.expect("escape_json の出力は JSON としてパースできること");
let value = String::try_from(parsed.value()).expect("JSON 文字列値として読めること");
assert_eq!(value, s);
}
}
#[test]
fn json_array_and_object_roundtrip_samples() {
let items = vec!["a".to_string(), "b\"c".to_string()];
let json = json_array(&items);
let parsed = nojson::RawJson::parse(&json).expect("json_array の出力はパースできること");
let arr = parsed.value().to_array().expect("配列であること");
let actual: Vec<String> = arr
.map(|item| String::try_from(item).expect("配列要素は文字列であること"))
.collect();
assert_eq!(actual, items);
let mut map = BTreeMap::new();
map.insert("k".to_string(), "v".to_string());
let json = json_object(&map);
let parsed = nojson::RawJson::parse(&json).expect("json_object の出力はパースできること");
let obj = parsed.value().to_object().expect("オブジェクトであること");
let mut actual = BTreeMap::new();
for (k, v) in obj {
actual.insert(
String::try_from(k).expect("キーは文字列であること"),
String::try_from(v).expect("値は文字列であること"),
);
}
assert_eq!(actual, map);
}
#[test]
fn percent_encode_encodes_slash_in_image_refs() {
assert_eq!(
percent_encode_path_segment("ghcr.io/org/app:1.0"),
"ghcr.io%2Forg%2Fapp%3A1.0"
);
assert_eq!(
percent_encode_component("ghcr.io/org/app"),
"ghcr.io%2Forg%2Fapp"
);
}
#[test]
fn parse_ports_extracts_tcp_udp_sctp() {
let json = r#"{"NetworkSettings":{"Ports":{"80/tcp":[{"HostPort":"8080"}],"53/udp":[{"HostPort":"5353"}],"5060/sctp":[{"HostPort":"5061"}]}}}"#;
let parsed = nojson::RawJson::parse(json).expect("処理に失敗しないこと");
let ports = parse_ports(&parsed);
assert_eq!(
ports.map_to_host_port_ipv4(80.tcp()),
Some(8080),
"tcp ポートが取得できること"
);
assert_eq!(
ports.map_to_host_port_ipv4(53.udp()),
Some(5353),
"udp ポートが取得できること"
);
assert_eq!(
ports.map_to_host_port_ipv4(5060.sctp()),
Some(5061),
"sctp ポートが取得できること"
);
}
#[test]
fn parse_ports_returns_empty_when_network_settings_missing() {
let json = r#"{}"#;
let parsed = nojson::RawJson::parse(json).expect("処理に失敗しないこと");
let ports = parse_ports(&parsed);
assert!(
ports.map_to_host_port_ipv4(80.tcp()).is_none(),
"NetworkSettings 欠落時は空の Ports になること"
);
}
#[test]
fn parse_ports_returns_empty_when_ports_missing() {
let json = r#"{"NetworkSettings":{}}"#;
let parsed = nojson::RawJson::parse(json).expect("処理に失敗しないこと");
let ports = parse_ports(&parsed);
assert!(
ports.map_to_host_port_ipv4(80.tcp()).is_none(),
"Ports 欠落時は空の Ports になること"
);
}
#[test]
fn parse_ports_ignores_invalid_port_strings() {
let json = r#"{"NetworkSettings":{"Ports":{"abc/tcp":[{"HostPort":"8080"}],"80/xyz":[{"HostPort":"8080"}]}}}"#;
let parsed = nojson::RawJson::parse(json).expect("処理に失敗しないこと");
let ports = parse_ports(&parsed);
assert!(
ports.map_to_host_port_ipv4(80.tcp()).is_none(),
"不正なポート文字列は無視されること"
);
}
#[test]
fn build_exposed_ports_formats_protocols() {
let ports = vec![
PortMapping {
container_port: 80.into(),
host_port: 8080,
},
PortMapping {
container_port: 53.udp(),
host_port: 5353,
},
PortMapping {
container_port: 5060.sctp(),
host_port: 5061,
},
];
let exposed = build_exposed_ports(&ports);
assert_eq!(
exposed,
vec!["80/tcp", "53/udp", "5060/sctp"],
"tcp/udp/sctp の書式が正しいこと"
);
}
fn config_with_mounts(mounts: Vec<crate::core::mounts::Mount>) -> ContainerConfig {
ContainerConfig {
image: "alpine:latest".into(),
entrypoint: None,
cmd: vec![],
env: vec![],
ports: vec![],
mounts,
name: None,
labels: BTreeMap::new(),
privileged: false,
working_dir: None,
user: None,
init: false,
health_check: None,
cap_add: vec![],
cap_drop: vec![],
shm_size: None,
readonly_rootfs: false,
hostname: None,
open_stdin: None,
network: None,
platform: None,
extra_hosts: vec![],
}
}
#[test]
fn from_config_bind_readonly_gets_ro_suffix() {
use crate::core::mounts::{AccessMode, Mount};
let mount = Mount::bind_mount("/host", "/container").with_access_mode(AccessMode::ReadOnly);
let body = CreateContainerBody::from_config(config_with_mounts(vec![mount]))
.expect("Bind マウントは成功すること");
assert_eq!(
body.host_config.binds,
vec!["/host:/container:ro".to_string()],
"ReadOnly は :ro サフィックスになること"
);
}
#[test]
fn from_config_bind_readwrite_gets_rw_suffix() {
use crate::core::mounts::{AccessMode, Mount};
let mount =
Mount::bind_mount("/host", "/container").with_access_mode(AccessMode::ReadWrite);
let body = CreateContainerBody::from_config(config_with_mounts(vec![mount]))
.expect("Bind マウントは成功すること");
assert_eq!(
body.host_config.binds,
vec!["/host:/container:rw".to_string()],
"ReadWrite は :rw サフィックスになること"
);
}
#[test]
fn from_config_volume_mount_is_reflected_in_mounts() {
use crate::core::mounts::Mount;
let body = CreateContainerBody::from_config(config_with_mounts(vec![Mount::volume_mount(
"data",
"/container/data",
)]))
.expect("Volume マウントが成功すること");
let json = body.to_json_string().expect("JSON 出力に失敗した");
assert!(
json.contains("\"Type\":\"volume\""),
"Mounts に volume タイプが含まれること: {json}"
);
assert!(
json.contains("\"Source\":\"data\""),
"Mounts にソース名が含まれること: {json}"
);
assert!(
json.contains("\"Target\":\"/container/data\""),
"Mounts にターゲットが含まれること: {json}"
);
}
#[test]
fn from_config_tmpfs_mount_is_reflected_in_mounts() {
use crate::core::mounts::Mount;
let body = CreateContainerBody::from_config(config_with_mounts(vec![
Mount::tmpfs_mount("/tmpfs")
.with_size_bytes(1_000_000)
.with_mode(0o1777),
]))
.expect("Tmpfs マウントが成功すること");
let json = body.to_json_string().expect("JSON 出力に失敗した");
assert!(
json.contains("\"Type\":\"tmpfs\""),
"Mounts に tmpfs タイプが含まれること: {json}"
);
assert!(
json.contains("\"Target\":\"/tmpfs\""),
"Mounts にターゲットが含まれること: {json}"
);
assert!(
json.contains("\"SizeBytes\":1000000"),
"TmpfsOptions にサイズが含まれること: {json}"
);
assert!(
json.contains("\"Mode\":1023"),
"TmpfsOptions にモードが含まれること: {json}"
);
}
#[test]
fn split_pull_reference_accepts_registry_port_and_digest() {
assert_eq!(
split_pull_reference("localhost:5000/nginx:latest"),
("localhost:5000/nginx", "latest")
);
assert_eq!(
split_pull_reference("localhost:5000/nginx"),
("localhost:5000/nginx", "latest")
);
assert_eq!(split_pull_reference("nginx:1.25"), ("nginx", "1.25"));
assert_eq!(split_pull_reference("nginx"), ("nginx", "latest"));
assert_eq!(
split_pull_reference("docker.io/library/nginx:latest"),
("docker.io/library/nginx", "latest")
);
assert_eq!(
split_pull_reference("nginx@sha256:ab12cd"),
("nginx", "sha256:ab12cd")
);
assert_eq!(
split_pull_reference("nginx:latest@sha256:ab12cd"),
("nginx", "sha256:ab12cd")
);
assert_eq!(
split_pull_reference("localhost:5000/nginx:1.25@sha256:ab12cd"),
("localhost:5000/nginx", "sha256:ab12cd")
);
assert_eq!(
split_pull_reference("nginx:sha256:ab12cd"),
("nginx", "sha256:ab12cd")
);
assert_eq!(split_pull_reference("nginx:"), ("nginx", "latest"));
assert_eq!(split_pull_reference("nginx@"), ("nginx", "latest"));
}
fn expect_other(err: ClientError) -> String {
match err {
ClientError::Other(msg) => msg,
other => panic!("Other 以外のエラー: {other}"),
}
}
fn expect_json(err: ClientError) {
match err {
ClientError::Json(_) => {}
other => panic!("Json 以外のエラー: {other}"),
}
}
#[test]
fn check_pull_stream_errors_detects_error_and_detail() {
let err = check_pull_stream_errors(
br#"{"error":"denied","errorDetail":{"message":"detail denied"}}"#,
)
.expect_err("error がある行は失敗すること");
assert_eq!(expect_other(err), "denied");
let err = check_pull_stream_errors(br#"{"error":"only error"}"#)
.expect_err("error のみでも失敗すること");
assert_eq!(expect_other(err), "only error");
let err = check_pull_stream_errors(br#"{"errorDetail":{"message":"only detail"}}"#)
.expect_err("errorDetail.message のみでも失敗すること");
assert_eq!(expect_other(err), "only detail");
let err = check_pull_stream_errors(br#"{"error":"","errorDetail":{"message":"real"}}"#)
.expect_err("空 error + 非空 detail は失敗すること");
assert_eq!(expect_other(err), "real");
}
#[test]
fn check_pull_stream_errors_ignores_empty_error_fields() {
check_pull_stream_errors(br#"{"error":""}"#).expect("空 error は成功であること");
check_pull_stream_errors(br#"{"errorDetail":{}}"#)
.expect("空 errorDetail は成功であること");
}
#[test]
fn check_pull_stream_errors_accepts_progress_and_aux() {
let body =
b"{\"status\":\"Pulling from library/nginx\"}\n{\"status\":\"Download complete\"}\n";
check_pull_stream_errors(body).expect("進捗のみは成功であること");
check_pull_stream_errors(br#"{"aux":{"ID":"sha256:dead"}}"#)
.expect("aux のみは成功であること");
}
#[test]
fn check_pull_stream_errors_rejects_non_object_and_bad_json() {
expect_json(check_pull_stream_errors(b"[]").expect_err("配列行は Json エラーであること"));
expect_json(
check_pull_stream_errors(b"{not-json").expect_err("不正 JSON は Json エラーであること"),
);
expect_json(
check_pull_stream_errors(&[0xff, 0xfe]).expect_err("非 UTF-8 は Json エラーであること"),
);
}
#[test]
fn check_pull_stream_errors_fails_on_error_after_progress() {
let body = b"{\"status\":\"Pulling\"}\n{\"error\":\"boom\"}\n";
let err = check_pull_stream_errors(body).expect_err("進捗後の error は失敗すること");
assert_eq!(expect_other(err), "boom");
}
#[test]
fn check_pull_stream_errors_accepts_empty_body() {
check_pull_stream_errors(b"").expect("空ボディは成功であること");
check_pull_stream_errors(b"\n\n").expect("空行のみは成功であること");
}
#[test]
fn encode_docker_api_request_sets_connection_close() {
let bytes =
encode_docker_api_request("GET", "/version", None).expect("エンコードに失敗しないこと");
let text = String::from_utf8(bytes).expect("UTF-8 であること");
assert!(
text.contains("Connection: close\r\n"),
"Connection: close が含まれること: {text}"
);
}
#[test]
fn read_http11_response_returns_before_peer_closes() {
use std::sync::mpsc;
use std::time::Duration;
let (mut writer, mut reader) =
UnixStream::pair().expect("UnixStream::pair に失敗しないこと");
writer
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhello")
.expect("レスポンス書き込みに失敗しないこと");
let (tx, rx) = mpsc::channel();
std::thread::spawn(move || {
let result = read_http11_response(&mut reader, "GET", BodyLimit::Unlimited);
let _ = tx.send(result);
});
let response = rx
.recv_timeout(Duration::from_millis(200))
.expect("ボディ完了後に keep-alive 待ちでハングしないこと")
.expect("デコードに失敗しないこと");
assert_eq!(
response.body_bytes(),
Some(b"hello".as_slice()),
"Content-Length ボディが取得できること"
);
drop(writer);
}
fn make_frame(stream_type: u8, payload: &[u8]) -> Vec<u8> {
let mut frame = vec![stream_type, 0, 0, 0];
frame.extend_from_slice(&(payload.len() as u32).to_be_bytes());
frame.extend_from_slice(payload);
frame
}
#[test]
fn demux_exec_stream_separates_stdout_and_stderr() {
let mut data = make_frame(1, b"hello ");
data.extend(make_frame(2, b"err "));
data.extend(make_frame(1, b"world"));
data.extend(make_frame(2, b"msg"));
let (stdout, stderr) = demux_exec_stream(&data);
assert_eq!(stdout, b"hello world", "stdout が正しく結合されること");
assert_eq!(stderr, b"err msg", "stderr が正しく結合されること");
}
#[test]
fn demux_exec_stream_handles_empty_input() {
let (stdout, stderr) = demux_exec_stream(b"");
assert!(stdout.is_empty(), "空入力では stdout が空であること");
assert!(stderr.is_empty(), "空入力では stderr が空であること");
}
#[test]
fn demux_exec_stream_ignores_stdin_and_unknown_types() {
let mut data = make_frame(0, b"stdin");
data.extend(make_frame(3, b"unknown"));
data.extend(make_frame(1, b"out"));
let (stdout, stderr) = demux_exec_stream(&data);
assert_eq!(stdout, b"out", "stdout のみ取得されること");
assert!(stderr.is_empty(), "stderr は空であること");
}
#[test]
fn demux_exec_stream_handles_zero_length_payload() {
let mut data = make_frame(1, b"");
data.extend(make_frame(1, b"data"));
let (stdout, stderr) = demux_exec_stream(&data);
assert_eq!(stdout, b"data", "payload_len=0 の後も正しく処理されること");
assert!(stderr.is_empty());
}
#[test]
fn demux_exec_stream_handles_truncated_frame() {
let mut data = make_frame(1, b"hello");
data.extend_from_slice(&[2, 0, 0, 0, 0, 0, 0, 10]);
data.extend_from_slice(b"par");
let (stdout, stderr) = demux_exec_stream(&data);
assert_eq!(stdout, b"hello", "完全なフレームは正しく処理されること");
assert_eq!(stderr, b"par", "切断フレームは部分出力を返すこと");
}
#[test]
fn classify_archive_404_path_not_found_message() {
let body = br#"{"message":"Could not find the file /no/such in container abc123"}"#;
let err = classify_archive_404("abc123", "/no/such", body);
assert!(
matches!(err, ClientError::ContainerPathNotFound(ref p) if p == "/no/such"),
"ContainerPathNotFound にパスが保持されること: {err:?}"
);
assert_eq!(err.to_string(), "container path not found: /no/such");
}
#[test]
fn classify_archive_404_container_not_found_message() {
let body = br#"{"message":"No such container: abc123"}"#;
let err = classify_archive_404("abc123", "/etc/hostname", body);
assert!(
matches!(err, ClientError::ContainerNotFound(ref id) if id == "abc123"),
"ContainerNotFound に ID が保持されること: {err:?}"
);
}
#[test]
fn classify_archive_404_falls_back_on_unknown_message() {
for body in [
&b""[..],
b"not-json",
br#"{"error":"something else"}"#,
br#"{"message":"an unknown daemon message"}"#,
] {
let err = classify_archive_404("abc123", "/etc/hostname", body);
assert!(
matches!(err, ClientError::ContainerNotFound(_)),
"フォールバックは ContainerNotFound であること: {body:?} → {err:?}"
);
}
}
#[test]
fn classify_archive_404_prefix_match_not_contains() {
let body =
br#"{"message":"Could not find the file /etc/No such container: x in container abc123"}"#;
let err = classify_archive_404("abc123", "/etc/No such container: x", body);
assert!(
matches!(err, ClientError::ContainerPathNotFound(_)),
"先頭一致で分類されること: {err:?}"
);
}
#[test]
fn parse_daemon_error_message_rejects_non_string_and_non_object() {
for body in [
&br#"{"message":123}"#[..],
b"[1,2,3]",
b"\"just a string\"",
&[0xff, 0xfe][..],
] {
assert!(
parse_daemon_error_message(body).is_none(),
"不正な message は None になること: {body:?}"
);
}
}
#[test]
fn parse_daemon_error_message_extracts_message() {
assert_eq!(
parse_daemon_error_message(br#"{"message":"No such container: abc"}"#).as_deref(),
Some("No such container: abc")
);
}
#[test]
fn resolve_exec_exit_code_picks_first_finished_state() {
let states = [(true, None), (true, None), (false, Some(3))];
assert_eq!(resolve_exec_exit_code(&states), Some(3));
}
#[test]
fn resolve_exec_exit_code_returns_none_when_always_running() {
let states = [(true, None), (true, None), (true, None)];
assert_eq!(resolve_exec_exit_code(&states), None);
}
#[test]
fn resolve_exec_exit_code_handles_first_state_finished() {
let states = [(false, Some(0))];
assert_eq!(resolve_exec_exit_code(&states), Some(0));
}
#[test]
fn exec_exit_code_backoff_sequence_matches_spec() {
assert_eq!(EXEC_EXIT_CODE_BACKOFF_MILLIS, &[10, 50, 200, 500, 1000]);
let attempts = EXEC_EXIT_CODE_BACKOFF_MILLIS.len() + 1;
let total_sleep: u64 = EXEC_EXIT_CODE_BACKOFF_MILLIS.iter().sum();
assert_eq!(
attempts, 6,
"試行回数が 6 回であること (初回即時 + 5 回の再試行)"
);
assert_eq!(total_sleep, 1760, "合計 sleep が約 1.76 秒であること");
}
#[test]
fn parse_exec_inspect_state_reads_running_and_exit_code() {
let (running, exit_code) = parse_exec_inspect_state(br#"{"Running":false,"ExitCode":7}"#)
.expect("正常応答はパースできること");
assert!(!running, "Running == false が読み取れること");
assert_eq!(exit_code, Some(7), "ExitCode が読み取れること");
let (running, _exit_code) = parse_exec_inspect_state(br#"{"Running":true,"ExitCode":0}"#)
.expect("正常応答はパースできること");
assert!(running, "Running == true が読み取れること");
}
#[test]
fn parse_exec_inspect_state_treats_missing_running_as_finished() {
let (running, exit_code) = parse_exec_inspect_state(br#"{"Running":"yes"}"#)
.expect("型不一致でも JSON として有効ならパースできること");
assert!(!running, "型不一致の Running は false 扱いであること");
assert_eq!(exit_code, None, "ExitCode が無い場合は None であること");
let (running, _) =
parse_exec_inspect_state(br#"{}"#).expect("空オブジェクトはパースできること");
assert!(!running, "Running 欠落は false 扱いであること");
}
#[test]
fn parse_exec_inspect_state_rejects_invalid_json() {
assert!(
parse_exec_inspect_state(b"\xff\xfe").is_err(),
"非 UTF-8 は Json エラーになること"
);
assert!(
parse_exec_inspect_state(b"not json").is_err(),
"JSON パース失敗は Json エラーになること"
);
}
#[test]
fn pull_error_message_includes_daemon_message() {
let msg = pull_error_message(
"ghcr.io/org/app:1.0",
401,
br#"{"message":"unauthorized: authentication required"}"#,
);
assert_eq!(
msg,
"failed to pull image ghcr.io/org/app:1.0: 401: unauthorized: authentication required",
"status と message が含まれること: {msg}"
);
let msg = pull_error_message(
"ghcr.io/org/app:1.0",
500,
br#"{"message":"registry is out of service"}"#,
);
assert_eq!(
msg, "failed to pull image ghcr.io/org/app:1.0: 500: registry is out of service",
"5xx でも message が含まれること: {msg}"
);
}
#[test]
fn pull_error_message_falls_back_without_message() {
let status = 500;
let expected = format!("failed to pull image img:latest: {status}");
assert_eq!(
pull_error_message("img:latest", status, &[]),
expected,
"ボディ無しは現行文言に落ちること"
);
assert_eq!(
pull_error_message("img:latest", status, b"not json"),
expected,
"非 JSON は現行文言に落ちること"
);
assert_eq!(
pull_error_message(
"img:latest",
status,
br#"{"errorDetail":{"message":"boom"}}"#
),
expected,
"トップレベルの message 欠落は現行文言に落ちること"
);
assert_eq!(
pull_error_message("img:latest", status, br#"{"message":""}"#),
expected,
"空 message は現行文言に落ちること"
);
assert_eq!(
pull_error_message("img:latest", status, br#"{"message":null}"#),
expected,
"null の message は現行文言に落ちること"
);
assert_eq!(
pull_error_message("img:latest", status, br#"{"message":42}"#),
expected,
"数値の message は現行文言に落ちること"
);
assert_eq!(
pull_error_message("img:latest", status, b"[1,2,3]"),
expected,
"配列ボディは現行文言に落ちること"
);
assert_eq!(
pull_error_message("img:latest", status, b"\xff\xfe"),
expected,
"非 UTF-8 ボディは現行文言に落ちること"
);
}
}