diff --git a/src/exec.rs b/src/exec.rs index 8d07508..f077329 100644 --- a/src/exec.rs +++ b/src/exec.rs @@ -1,9 +1,9 @@ +use bollard::Docker; use bollard::exec::{CreateExecOptions, ResizeExecOptions, StartExecResults}; use bollard::query_parameters::LogsOptions; -use bollard::Docker; use futures_util::StreamExt; use miette::{IntoDiagnostic, Report, Result, WrapErr}; -use std::io::{stdout, Read, Stdin, Write}; +use std::io::{Read, Stdin, Write, stdout}; use std::time::Duration; use termion::raw::IntoRawMode; use termion::{async_stdin, is_tty, terminal_size}; @@ -206,6 +206,7 @@ pub async fn container_logs( Ok(()) } +#[derive(Debug, Clone, Copy)] pub struct ExitCode(i64); impl ExitCode { @@ -213,7 +214,7 @@ impl ExitCode { self.0 == 0 } - pub fn to_result(&self) -> Result<()> { + pub fn to_result(self) -> Result<()> { match self.0 { 0 => Ok(()), code => Err(Report::msg(format!( diff --git a/src/main.rs b/src/main.rs index ffda7ec..bcfd7e5 100644 --- a/src/main.rs +++ b/src/main.rs @@ -16,13 +16,18 @@ use crate::service::ServiceTrait; use crate::service::{RedisTls, Service}; use crate::sources::Sources; use bollard::Docker; +use futures_util::future::{Either, join_all, select}; use itertools::Itertools; use miette::{IntoDiagnostic, Report, Result, WrapErr}; +use std::collections::HashMap; use std::env::{var, vars}; use std::fs::{create_dir_all, remove_dir_all, remove_file, write}; use std::io::stdout; use std::os::unix::process::CommandExt; +use std::pin::pin; use std::process::{Command, ExitCode}; +use std::time::Duration; +use tokio::time::sleep; mod args; mod cloud; @@ -92,23 +97,24 @@ async fn main() -> Result { let cache_dir = cache_dir.into_diagnostic()?; if let Some(id) = cache_dir.file_name().to_str() && id.starts_with("haze-") - && !retain.iter().any(|cloud| cloud.id == id) { - let path = cache_dir.path(); - // By default, the worktrees live within the work dir, we do not want to delete it here. - if path == config.worktree_dir { - continue; - } - - if path.is_dir() { - remove_dir_all(&path).into_diagnostic().wrap_err_with(|| { - format!("Failed to cleanup stray cache dir: {}", path.display()) - })?; - } else { - remove_file(&path).into_diagnostic().wrap_err_with(|| { - format!("Failed to cleanup stray cache file: {}", path.display()) - })?; - } + && !retain.iter().any(|cloud| cloud.id == id) + { + let path = cache_dir.path(); + // By default, the worktrees live within the work dir, we do not want to delete it here. + if path == config.worktree_dir { + continue; } + + if path.is_dir() { + remove_dir_all(&path).into_diagnostic().wrap_err_with(|| { + format!("Failed to cleanup stray cache dir: {}", path.display()) + })?; + } else { + remove_file(&path).into_diagnostic().wrap_err_with(|| { + format!("Failed to cleanup stray cache file: {}", path.display()) + })?; + } + } } prune_worktrees(&config)?; @@ -146,7 +152,7 @@ async fn main() -> Result { } } HazeArgs::Start { options } => { - setup(&docker, options, &config).await?; + setup(&docker, options, &config, false, false).await?; } HazeArgs::Stop { filter } => { let cloud = Cloud::get_by_filter(&docker, filter, &config).await?; @@ -277,37 +283,7 @@ async fn main() -> Result { } } HazeArgs::Test { options, mut args } => { - let cloud = Cloud::create(&docker, options, &config).await?; - println!("Waiting for servers to start"); - cloud.wait_for_start(&docker).await?; - - if !cloud.preset_config.is_empty() { - println!("Writing preset config"); - let encoded_preset_config = - serde_json::to_string(&cloud.preset_config).into_diagnostic()?; - cloud - .write_file(&docker, "config/preset.config.json", encoded_preset_config) - .await?; - cloud.write_file(&docker, "config/preset.config.php", "::default(), - ) - .await - { - cloud.destroy(&docker, &config).await?; - return Err(e); - } + let cloud = setup(&docker, options, &config, true, true).await?; if let Some(app) = args .first() .as_ref() @@ -326,7 +302,7 @@ async fn main() -> Result { return Ok(result.into()); } HazeArgs::Integration { options, mut args } => { - let cloud = setup(&docker, options, &config).await?; + let cloud = setup(&docker, options, &config, true, true).await?; args.insert(0, "integration".to_string()); cloud.exec(&docker, args, false, get_forward_env()).await?; @@ -398,7 +374,7 @@ async fn main() -> Result { cloud.destroy(&docker, &config).await?; } HazeArgs::Shell { command, options } => { - let cloud = setup(&docker, options, &config).await?; + let cloud = setup(&docker, options, &config, true, false).await?; cloud .exec( &docker, @@ -566,11 +542,17 @@ async fn main() -> Result { Ok(ExitCode::SUCCESS) } -async fn setup(docker: &Docker, options: CloudOptions, config: &HazeConfig) -> Result { +async fn setup( + docker: &Docker, + options: CloudOptions, + config: &HazeConfig, + always_setup: bool, + force_wait: bool, +) -> Result { let cloud = Cloud::create(docker, options, config).await?; println!("{}", cloud.address); let host = cloud.address.split_once("://").expect("no address?").1; - if config.auto_setup.enabled { + if always_setup || config.auto_setup.enabled { println!("Waiting for servers to start"); cloud.wait_for_start(docker).await?; @@ -703,5 +685,37 @@ async fn setup(docker: &Docker, options: CloudOptions, config: &HazeConfig) -> R .await?; } } + + { + let mut waiting = HashMap::new(); + for service in cloud.services() { + // give the service 100ms before we mark them as pending + let sleep = pin!(sleep(Duration::from_millis(100))); + let future = service.final_wait(docker, &cloud.id); + if let Either::Left((_, future)) = select(sleep, future).await { + waiting.insert(service.name(), future); + } + } + + if !waiting.is_empty() { + println!(); + println!("The folowing services are still starting in the background:"); + for service in waiting.keys() { + println!(" - {service}"); + } + println!(); + if !force_wait { + println!( + "You can safely terminate this command with CTRL-C and use the instance, but some functionality provided by the service might not be ready yet." + ) + } + + join_all(waiting.into_values()) + .await + .into_iter() + .collect::, _>>()?; + } + } + Ok(cloud) } diff --git a/src/script.rs b/src/script.rs index fe4d026..4154be0 100644 --- a/src/script.rs +++ b/src/script.rs @@ -57,7 +57,7 @@ pub async fn run_script( .dont_create(), ); - let cloud = setup(docker, options, config).await?; + let cloud = setup(docker, options, config, true, true).await?; cloud .exec( docker, diff --git a/src/service.rs b/src/service.rs index a70413e..233e88f 100644 --- a/src/service.rs +++ b/src/service.rs @@ -205,6 +205,14 @@ pub trait ServiceTrait { fn exec_shell(&self) -> &'static str { "bash" } + + /// Allow services to delay the completion of the `start` command untill they are fully setup + /// + /// This is not for services which are critical for the running of the instance (use `is_healthy` for that) + /// But instead for cases where the instance is mostly usable, but might have some missing features + async fn final_wait(&self, _docker: &Docker, _cloud_id: &str) -> Result<()> { + Ok(()) + } } #[derive(Clone, Eq, PartialEq, Debug)] diff --git a/src/service/authentik/mod.rs b/src/service/authentik/mod.rs index dd7b364..86714f1 100644 --- a/src/service/authentik/mod.rs +++ b/src/service/authentik/mod.rs @@ -2,13 +2,16 @@ mod oidc; mod saml; mod scim; +use futures_util::future::join; pub use oidc::AuthentikOidc; pub use saml::AuthentikSaml; pub use scim::AuthentikScim; +use tokio::time::{sleep, timeout}; use crate::Result; use crate::cloud::CloudOptions; use crate::config::{HazeConfig, ProxyConfig}; +use crate::exec::exec; use crate::image::pull_image; use crate::service::ServiceTrait; use bollard::Docker; @@ -19,6 +22,7 @@ use maplit::hashmap; use miette::{IntoDiagnostic, Report, WrapErr}; use std::fs::{create_dir_all, write}; use std::net::{IpAddr, Ipv4Addr}; +use std::time::Duration; const AUTHENTIK_IMAGE: &str = "ghcr.io/goauthentik/server:2026.8.0"; const POSTGRES_IMAGE: &str = "docker.io/library/postgres:16-alpine"; @@ -308,4 +312,75 @@ Authentik groups: 'authentik-group1', 'authentik-group2' and 'authentik-group3' "#, ))) } + + async fn final_wait(&self, docker: &Docker, cloud_id: &str) -> Result<()> { + let server = timeout(Duration::from_mins(15), wait_server(docker, cloud_id)); + let worker = timeout(Duration::from_mins(15), wait_worker(docker, cloud_id)); + + let (server_result, worker_result) = join(server, worker).await; + match server_result { + Err(_) => { + eprintln!("timeout while waiting for authentic server to become ready"); + } + Ok(res) => res?, + } + match worker_result { + Err(_) => { + eprintln!("timeout while waiting for authentic worker to become ready"); + } + Ok(res) => res?, + } + + Ok(()) + } +} + +async fn wait_server(docker: &Docker, cloud_id: &str) -> Result<()> { + loop { + sleep(Duration::from_secs(1)).await; + if healthcheck_server(docker, cloud_id).await? { + break; + } + } + Ok(()) +} + +async fn wait_worker(docker: &Docker, cloud_id: &str) -> Result<()> { + loop { + sleep(Duration::from_secs(1)).await; + if healthcheck_worker(docker, cloud_id).await? { + break; + } + } + Ok(()) +} + +async fn healthcheck_server(docker: &Docker, cloud_id: &str) -> Result { + let mut output = Vec::new(); + let exit = exec( + docker, + cloud_id, + "root", + vec!["curl", "http://authentik:9000/-/health/ready"], + Vec::::default(), + Some(&mut output), + ) + .await?; + Ok(exit.is_ok()) +} + +async fn healthcheck_worker(docker: &Docker, cloud_id: &str) -> Result { + let container = container_name(cloud_id, "worker"); + + let mut output = Vec::new(); + let exit = exec( + docker, + container, + "root", + vec!["ak", "healthcheck"], + Vec::::default(), + Some(&mut output), + ) + .await?; + Ok(exit.is_ok()) }