use crate::compose::types::ComposeFile;
use crate::error::{ComposeError, Result};
use super::filter_services;
use super::parallel::{
filter_levels, first_error, join_bounded, restart_service_set, retain_levels,
};
use super::targets::{stop_deadline, stop_timeout_param};
use crate::engine::Engine;
use crate::libpod::API_PREFIX;
const WAIT_NAME_WIDTH: usize = 32;
fn wait_header() -> String {
format!("{:<WAIT_NAME_WIDTH$} EXIT", "NAME")
}
fn wait_row(container: &str, code: i64) -> String {
let cell = crate::ui::fit_cell(container, WAIT_NAME_WIDTH);
let style = if code == 0 {
crate::ui::Style::new().dimmed()
} else {
crate::ui::Style::new().fg_color(Some(crate::ui::AnsiColor::Red.into()))
};
let coloured = crate::ui::stdout_colored();
format!(
"{} {}",
crate::ui::paint(crate::ui::identity_style(container), &cell, coloured),
crate::ui::paint(style, &code.to_string(), coloured)
)
}
fn note_if_idle(acted: &std::sync::atomic::AtomicBool, verb: &str) {
if !acted.load(std::sync::atomic::Ordering::Relaxed) {
crate::ui::progress_note(&format!("no containers to {verb}"));
}
}
impl Engine {
pub(super) async fn run_lifecycle_op(
&self,
path: &str,
container: &str,
done: &str,
) -> Result<bool> {
match self.client.post_empty_ok(path).await {
Ok(()) => {
crate::ui::progress_line("Container", container, done);
Ok(true)
}
Err(e) if e.is_status(304) || e.is_status(404) || e.is_kill_of_stopped() => {
tracing::debug!("{container}: {done} skipped ({e})");
Ok(false)
}
Err(e) => Err(ComposeError::Podman(e)),
}
}
pub(super) async fn run_idempotent_state_op(
&self,
path: &str,
container: &str,
done: &str,
) -> Result<bool> {
match self.client.post_empty_ok(path).await {
Ok(()) => {
crate::ui::progress_line("Container", container, done);
Ok(true)
}
Err(e) if e.is_status(304) || e.is_status(404) || e.is_state_conflict() => {
tracing::debug!("{container}: {done} skipped ({e})");
Ok(false)
}
Err(e) => Err(ComposeError::Podman(e)),
}
}
pub(super) async fn stop_container(&self, container: &str, grace: i32) -> Result<()> {
let path = format!(
"{API_PREFIX}/containers/{}/stop?t={}",
crate::libpod::urlencoded(container),
stop_timeout_param(grace),
);
match self
.client
.post_empty_ok_within(&path, stop_deadline(grace))
.await
{
Ok(()) => {
crate::ui::progress_line("Container", container, "Stopped");
Ok(())
}
Err(e) if e.is_status(304) || e.is_status(404) => {
tracing::debug!("{container}: stop skipped ({e})");
Ok(())
}
Err(e) if e.is_timeout() => {
tracing::warn!(
"{container}: stop did not complete within the grace window; escalating to SIGKILL"
);
let kill_path = format!(
"{API_PREFIX}/containers/{}/kill?signal=SIGKILL",
crate::libpod::urlencoded(container),
);
match self.client.post_empty_ok(&kill_path).await {
Ok(()) => {
crate::ui::progress_line(
"Container",
container,
"Killed (after stop timeout)",
);
Ok(())
}
Err(e) if e.is_status(404) || e.is_status(409) => {
tracing::debug!("{container}: SIGKILL skipped ({e})");
Ok(())
}
Err(e) => Err(ComposeError::Podman(e)),
}
}
Err(e) => Err(ComposeError::Podman(e)),
}
}
pub async fn restart(&self, file: &ComposeFile, service_name: Option<&str>) -> Result<()> {
let targets: Vec<String> = service_name
.map(|s| vec![s.to_string()])
.unwrap_or_default();
self.restart_with_options(file, &targets, false).await
}
pub async fn restart_with_options(
&self,
file: &ComposeFile,
target_services: &[String],
no_deps: bool,
) -> Result<()> {
super::targets::validate_targets(file, target_services)?;
let (restart_set, targets) = restart_service_set(file, target_services, no_deps);
let levels = retain_levels(crate::compose::resolve_levels(file)?, |n| {
restart_set.contains(n)
});
let acted = std::sync::atomic::AtomicBool::new(false);
let mut first_err: Option<ComposeError> = None;
for level in &levels {
let futs = level.iter().map(|name| {
let service = &file.services[name];
let done = if targets.contains(name) {
"Restarted"
} else {
"Restarted (dependency)"
};
self.restart_one_service(name, service, done, &acted)
});
if let Some(e) = first_error(join_bounded(futs).await) {
first_err.get_or_insert(e);
}
}
if let Some(e) = first_err {
return Err(e);
}
note_if_idle(&acted, "restart");
Ok(())
}
pub async fn wait_services(
&self,
file: &ComposeFile,
target_services: &[String],
) -> Result<()> {
self.wait_services_with_options(file, target_services, false)
.await
}
pub async fn wait_services_with_options(
&self,
file: &ComposeFile,
target_services: &[String],
json: bool,
) -> Result<()> {
let order = if target_services.is_empty() {
let order = crate::compose::resolve_order(file)?;
filter_services(file, order, &[])?
} else {
for name in target_services {
if !file.services.contains_key(name) {
return Err(ComposeError::ServiceNotFound(name.clone()));
}
}
let mut seen = std::collections::HashSet::new();
target_services
.iter()
.filter(|n| seen.insert(n.as_str()))
.cloned()
.collect::<Vec<_>>()
};
let mut last_nonzero = 0i64;
let mut printed_header = false;
for name in &order {
for container_name in self
.list_project_container_names(Some(name.as_str()))
.await?
{
let path = format!(
"{API_PREFIX}/containers/{}/wait?condition=stopped",
crate::libpod::urlencoded(&container_name),
);
let code = self
.client
.post_empty_json_unbounded::<i64>(&path)
.await
.map_err(ComposeError::Podman)?;
if code != 0 {
last_nonzero = code;
}
if json {
println!(
"{}",
serde_json::json!({ "Container": container_name, "ExitCode": code })
);
} else {
if !printed_header {
crate::ui::print_bold_header(&wait_header());
printed_header = true;
}
println!("{}", wait_row(&container_name, code));
}
}
}
if last_nonzero != 0 {
return Err(ComposeError::RunExited(last_nonzero));
}
Ok(())
}
pub async fn stop(&self, file: &ComposeFile, target_services: &[String]) -> Result<()> {
let mut levels = crate::compose::resolve_levels(file)?;
levels.reverse();
let levels = filter_levels(file, levels, target_services)?;
let acted = std::sync::atomic::AtomicBool::new(false);
for level in &levels {
let futs = level.iter().map(|name| {
let grace = self.grace_period_secs(&file.services[name]);
self.stop_one_service(name, grace, &acted)
});
if let Some(e) = first_error(join_bounded(futs).await) {
return Err(e);
}
}
note_if_idle(&acted, "stop");
Ok(())
}
pub async fn start(&self, file: &ComposeFile, target_services: &[String]) -> Result<()> {
let levels = crate::compose::resolve_levels(file)?;
let levels = filter_levels(file, levels, target_services)?;
let any_live = std::sync::atomic::AtomicBool::new(false);
let mut first_err: Option<ComposeError> = None;
for level in &levels {
let futs = level
.iter()
.map(|name| self.start_one_service(name, &any_live));
if let Some(e) = first_error(join_bounded(futs).await) {
first_err.get_or_insert(e);
}
}
if let Some(e) = first_err {
return Err(e);
}
note_if_idle(&any_live, "start (project not created)");
Ok(())
}
pub async fn kill(
&self,
file: &ComposeFile,
target_services: &[String],
signal: &str,
) -> Result<()> {
super::signal::validate_signal(signal)?;
let levels = crate::compose::resolve_levels(file)?;
let levels = filter_levels(file, levels, target_services)?;
let acted = std::sync::atomic::AtomicBool::new(false);
for level in &levels {
let futs = level
.iter()
.map(|name| self.kill_one_service(name, &file.services[name], signal, &acted));
if let Some(e) = first_error(join_bounded(futs).await) {
return Err(e);
}
}
note_if_idle(&acted, "signal");
Ok(())
}
pub async fn rm(
&self,
file: &ComposeFile,
target_services: &[String],
force: bool,
) -> Result<()> {
self.rm_with_options(file, target_services, force, false)
.await
}
pub async fn rm_with_options(
&self,
file: &ComposeFile,
target_services: &[String],
force: bool,
remove_volumes: bool,
) -> Result<()> {
let mut levels = crate::compose::resolve_levels(file)?;
levels.reverse();
let levels = filter_levels(file, levels, target_services)?;
let acted = std::sync::atomic::AtomicBool::new(false);
let mut first_err: Option<ComposeError> = None;
for level in &levels {
let futs = level.iter().map(|name| {
self.rm_one_service(name, &file.services[name], force, remove_volumes, &acted)
});
if let Some(e) = first_error(join_bounded(futs).await) {
first_err.get_or_insert(e);
}
}
if let Some(e) = first_err {
return Err(e);
}
note_if_idle(&acted, "remove");
Ok(())
}
pub async fn pause(&self, file: &ComposeFile, target_services: &[String]) -> Result<()> {
let levels = crate::compose::resolve_levels(file)?;
let levels = filter_levels(file, levels, target_services)?;
let acted = std::sync::atomic::AtomicBool::new(false);
let mut first_err: Option<ComposeError> = None;
for level in &levels {
let futs = level.iter().map(|name| {
self.idempotent_state_service(name, &file.services[name], "pause", "Paused", &acted)
});
if let Some(e) = first_error(join_bounded(futs).await) {
first_err.get_or_insert(e);
}
}
if let Some(e) = first_err {
return Err(e);
}
note_if_idle(&acted, "pause");
Ok(())
}
pub async fn unpause(&self, file: &ComposeFile, target_services: &[String]) -> Result<()> {
let levels = crate::compose::resolve_levels(file)?;
let levels = filter_levels(file, levels, target_services)?;
let acted = std::sync::atomic::AtomicBool::new(false);
let mut first_err: Option<ComposeError> = None;
for level in &levels {
let futs = level.iter().map(|name| {
self.idempotent_state_service(
name,
&file.services[name],
"unpause",
"Unpaused",
&acted,
)
});
if let Some(e) = first_error(join_bounded(futs).await) {
first_err.get_or_insert(e);
}
}
if let Some(e) = first_err {
return Err(e);
}
note_if_idle(&acted, "unpause");
Ok(())
}
pub(super) async fn container_exists(&self, name: &str) -> Result<bool> {
let path = format!(
"{API_PREFIX}/containers/{}/json",
crate::libpod::urlencoded(name),
);
match self.client.get_json::<serde_json::Value>(&path).await {
Ok(_) => Ok(true),
Err(e) if e.is_status(404) => Ok(false),
Err(e) => Err(ComposeError::Podman(e)),
}
}
}
#[cfg(test)]
mod wait_output_tests {
use super::{wait_header, wait_row, WAIT_NAME_WIDTH};
#[test]
fn the_exit_column_is_in_one_place() {
let code_col = |line: &str| {
let plain: String = strip_ansi(line);
plain.rfind(' ').map(|i| i + 1)
};
let header = wait_header();
let short = wait_row("a", 0);
let long = wait_row("project-service-name-12", 7);
assert_eq!(code_col(&header), code_col(&short));
assert_eq!(code_col(&header), code_col(&long));
assert_eq!(code_col(&header), Some(WAIT_NAME_WIDTH + 1));
}
#[test]
fn an_over_long_name_truncates() {
let row = strip_ansi(&wait_row(&"x".repeat(WAIT_NAME_WIDTH + 20), 0));
assert_eq!(row.chars().count(), WAIT_NAME_WIDTH + 2);
}
#[test]
fn a_container_name_cannot_drive_the_terminal() {
let row = wait_row("evil\u{1b}[31m\u{7}name", 0);
assert!(!row.contains('\u{1b}'), "{row:?}");
assert!(!row.contains('\u{7}'), "{row:?}");
assert!(row.contains("name"), "{row:?}");
}
fn strip_ansi(s: &str) -> String {
let mut out = String::new();
let mut chars = s.chars();
while let Some(c) = chars.next() {
if c == '\u{1b}' {
for c in chars.by_ref() {
if c == 'm' {
break;
}
}
} else {
out.push(c);
}
}
out
}
}