Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 56 additions & 0 deletions crates/clickhousectl/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -254,6 +254,62 @@ pub enum Error {
#[error("Server '{0}' is running; stop it first with `clickhousectl local server stop {0}`")]
ServerRunningCannotRemove(String),

#[error(
"Permission denied reading server metadata '{}': {source}. Check ownership and file permissions, then retry.",
path.display()
)]
ServerMetadataPermission {
path: PathBuf,
#[source]
source: std::io::Error,
},

#[error(
"Failed to read server metadata '{}': {source}. Check that the file is readable, then retry.",
path.display()
)]
ServerMetadataRead {
path: PathBuf,
#[source]
source: std::io::Error,
},

#[error(
"Server metadata '{}' is not valid UTF-8: {source}. Repair or remove the metadata file, then retry.",
path.display()
)]
ServerMetadataUtf8 {
path: PathBuf,
#[source]
source: std::string::FromUtf8Error,
},

#[error(
"Server metadata '{}' is not valid JSON: {source}. Repair or remove the metadata file, then retry.",
path.display()
)]
ServerMetadataParse {
path: PathBuf,
#[source]
source: serde_json::Error,
},

#[error("Failed to durably update server metadata '{}': {source}", path.display())]
ServerMetadataWrite {
path: PathBuf,
#[source]
source: std::io::Error,
},

#[error("Could not {operation} at '{}': {source}. {remediation}", path.display())]
ServerLock {
operation: &'static str,
path: PathBuf,
remediation: &'static str,
#[source]
source: std::io::Error,
},

#[error("{0}")]
Cloud(String),

Expand Down
19 changes: 11 additions & 8 deletions crates/clickhousectl/src/local/docker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1085,13 +1085,16 @@ pub async fn start_existing(docker: &Docker, id: &str) -> Result<()> {
/// Safe to call multiple times in one CLI invocation. When Docker isn't
/// reachable, `connect()` fails fast (no socket → immediate error, no I/O
/// timeout) and we return silently.
pub fn recover_project_postgres_blocking(project_cwd: &str) {
pub fn recover_project_postgres_blocking(
project_cwd: &str,
lock: &crate::local::server::MetadataLock,
) -> Result<()> {
use crate::local::server::{
Engine, ServerInfo, ensure_pg_data_dir, pg_instance_key, save_server_info,
server_meta_path_for_recovery,
Engine, ServerInfo, ensure_pg_data_dir, load_info_locked, pg_instance_key,
save_server_info_locked,
};
let cwd_owned = project_cwd.to_string();
let _ = block_on(async move {
block_on(async move {
let docker = match connect().await {
Ok(d) => d,
Err(_) => return Ok::<(), Error>(()),
Expand All @@ -1102,10 +1105,10 @@ pub fn recover_project_postgres_blocking(project_cwd: &str) {
};
for c in containers {
let key = pg_instance_key(&c.user_name, &c.major);
if server_meta_path_for_recovery(&key).exists() {
if load_info_locked(&key, lock)?.is_some() {
continue;
}
let _ = ensure_pg_data_dir(&c.user_name, &c.major);
ensure_pg_data_dir(&c.user_name, &c.major)?;
let info = ServerInfo {
name: key,
pid: 0,
Expand All @@ -1117,10 +1120,10 @@ pub fn recover_project_postgres_blocking(project_cwd: &str) {
engine: Engine::Postgres,
container_id: Some(c.container_id.clone()),
};
let _ = save_server_info(&info);
save_server_info_locked(&info, lock)?;
Comment thread
cursor[bot] marked this conversation as resolved.
}
Ok(())
});
})
}

#[cfg(test)]
Expand Down
87 changes: 60 additions & 27 deletions crates/clickhousectl/src/local/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -190,8 +190,8 @@ fn remove(version: &str, force: bool, json: bool) -> Result<()> {
// Recover orphaned servers so we detect a running process even when its
// metadata file is missing, then refuse to pull the binary out from under
// a server running on this version.
server::recover_current_project_servers();
let in_use: Vec<String> = server::list_running_servers()
server::recover_current_project_servers()?;
let in_use: Vec<String> = server::list_running_servers()?
.into_iter()
.filter(|i| i.version == version)
.map(|i| i.name)
Expand Down Expand Up @@ -274,19 +274,17 @@ fn run_client(
(h, p, v)
} else {
let server_name = name.as_deref().unwrap_or("default");
let entries = server::list_all_servers();
let entry = entries
.iter()
.find(|e| e.name == server_name)
let metadata_lock = server::lock_metadata()?;
server::recover_current_project_servers_locked(&metadata_lock)?;
let entry = server::server_entry_locked(server_name, &metadata_lock)?
.ok_or_else(|| Error::ServerNotFound(server_name.to_string()))?;
if !entry.running {
return Err(Error::ServerNotRunning(server_name.to_string()));
}
let info = entry
.info
.as_ref()
.ok_or_else(|| Error::ServerNotRunning(server_name.to_string()))?;
("localhost".to_string(), info.tcp_port, info.version.clone())
("localhost".to_string(), info.tcp_port, info.version)
};

let binary = paths::binary_path(&version)?;
Expand Down Expand Up @@ -358,6 +356,21 @@ fn resolve_direct_client_version(version_spec: Option<ClientVersionArg>) -> Resu
}
}

fn clean_up_untracked_child(child: &mut std::process::Child, primary: Error) -> Error {
let cleanup = match child.try_wait() {
Ok(Some(_)) => Ok(()),
Ok(None) => child.kill().and_then(|()| child.wait()).map(|_| ()),
Err(error) => Err(error),
};
match cleanup {
Ok(()) => primary,
Err(error) => Error::Exec(format!(
"{primary}; additionally failed to stop untracked PID {}: {error}",
child.id()
)),
}
}

#[allow(clippy::too_many_arguments)]
async fn start_server(
name: Option<String>,
Expand All @@ -375,12 +388,12 @@ async fn start_server(

// Recover any orphaned servers so name resolution and collision checks
// see processes that lost their metadata files.
server::recover_current_project_servers();
server::recover_current_project_servers()?;

// Resolve server name and check for collisions before any downloads
let server_name = server::resolve_name(name.as_deref())?;

if name.is_some() && server::is_server_running(&server_name) {
if name.is_some() && server::is_server_running(&server_name)? {
return Err(Error::ServerAlreadyRunning(server_name));
}

Expand Down Expand Up @@ -416,8 +429,18 @@ async fn start_server(
return Err(Error::VersionNotFound(version));
}

// Metadata locking is always acquired after version installation locks.
// Hold it from the final collision check through the process metadata
// commit so concurrent start/stop/client commands see one lifecycle state.
let metadata_lock = server::lock_metadata()?;
server::recover_current_project_servers_locked(&metadata_lock)?;
let server_name = server::resolve_name_locked(name.as_deref(), &metadata_lock)?;
if name.is_some() && server::is_server_running_locked(&server_name, &metadata_lock)? {
return Err(Error::ServerAlreadyRunning(server_name));
}

// Show running server count
let running = server::running_server_count();
let running = server::advisory_running_server_count_locked(&metadata_lock);
if running > 0 {
eprintln!(
"Note: {} server{} already running (use `clickhousectl local server list` to see them)",
Expand Down Expand Up @@ -509,7 +532,10 @@ async fn start_server(
engine: server::Engine::Clickhouse,
container_id: None,
};
server::save_server_info(&info)?;
if let Err(error) = server::save_server_info_locked(&info, &metadata_lock) {
return Err(clean_up_untracked_child(&mut child, error));
}
Comment thread
cursor[bot] marked this conversation as resolved.
drop(metadata_lock);

if no_wait {
server::check_spawn_health(&mut child, &server_name, &log_path).await?;
Expand Down Expand Up @@ -549,7 +575,10 @@ async fn start_server(
engine: server::Engine::Clickhouse,
container_id: None,
};
server::save_server_info(&info)?;
if let Err(error) = server::save_server_info_locked(&info, &metadata_lock) {
return Err(clean_up_untracked_child(&mut child, error));
}
drop(metadata_lock);

eprintln!(
"Server '{}' running (PID: {}, HTTP: {}, TCP: {})",
Expand Down Expand Up @@ -587,17 +616,15 @@ fn dotenv_server(
json: bool,
) -> Result<()> {
let server_name = name.unwrap_or("default");
let entries = server::list_all_servers();
let entry = entries
.iter()
.find(|e| e.name == server_name)
let metadata_lock = server::lock_metadata()?;
server::recover_current_project_servers_locked(&metadata_lock)?;
let entry = server::server_entry_locked(server_name, &metadata_lock)?
.ok_or_else(|| Error::ServerNotFound(server_name.to_string()))?;
if !entry.running {
return Err(Error::ServerNotRunning(server_name.to_string()));
}
let info = entry
.info
.as_ref()
.ok_or_else(|| Error::ServerNotRunning(server_name.to_string()))?;

// Only write vars we actually know from server metadata.
Expand Down Expand Up @@ -769,17 +796,18 @@ async fn run_server_commands(command: ServerCommands, json: bool) -> Result<()>

// Recover orphaned servers so we can stop processes
// that lost their metadata files.
server::recover_current_project_servers();
let metadata_lock = server::lock_metadata()?;
server::recover_current_project_servers_locked(&metadata_lock)?;

match classify_stop(
server::is_server_running(&name),
server::is_server_running_locked(&name, &metadata_lock)?,
server::server_data_dir(&name).exists(),
) {
StopOutcome::Stop => {
if !json {
println!("Stopping server '{}'...", name);
}
server::kill_server(&name)?;
server::kill_server_locked(&name, &metadata_lock)?;
let out = output::ServerStopOutput {
name,
already_stopped: false,
Expand Down Expand Up @@ -822,9 +850,10 @@ async fn run_server_commands(command: ServerCommands, json: bool) -> Result<()>

// Recover orphaned servers so we correctly detect a running
// process even when its metadata file is missing.
server::recover_current_project_servers();
let metadata_lock = server::lock_metadata()?;
server::recover_current_project_servers_locked(&metadata_lock)?;

if server::is_server_running(&name) {
if server::is_server_running_locked(&name, &metadata_lock)? {
return Err(Error::ServerRunningCannotRemove(name));
}
let data_dir = server::server_data_dir(&name);
Expand All @@ -834,7 +863,7 @@ async fn run_server_commands(command: ServerCommands, json: bool) -> Result<()>
// Remove the whole server directory (parent of data/)
let server_dir = data_dir.parent().unwrap();
std::fs::remove_dir_all(server_dir)?;
server::remove_server_info(&name);
server::try_remove_server_info_locked(&name, &metadata_lock)?;
let out = output::ServerRemoveOutput { name };
output::print_output(&out, json);
Ok(())
Expand Down Expand Up @@ -863,7 +892,7 @@ fn classify_stop(running: bool, exists_on_disk: bool) -> StopOutcome {
}

fn list_servers_local(json: bool) -> Result<()> {
let entries = server::list_all_servers();
let entries = server::list_all_servers()?;
let running_count = entries.iter().filter(|e| e.running).count();
let total = entries.len();

Expand Down Expand Up @@ -1006,8 +1035,12 @@ fn stop_server_global(name: &str, project: Option<&str>, json: bool) -> Result<(
}

fn stop_all_servers_local(json: bool) -> Result<()> {
let servers = server::list_running_servers();
let out = stop_servers(&servers, json, server::kill_server);
let metadata_lock = server::lock_metadata()?;
server::recover_current_project_servers_locked(&metadata_lock)?;
let servers = server::list_running_servers_locked(&metadata_lock)?;
let out = stop_servers(&servers, json, |name| {
server::kill_server_locked(name, &metadata_lock)
});
if json {
output::print_output(&out, json);
} else if servers.is_empty() {
Expand Down
Loading