use crate::error::{Error, Result};
use crate::init;
use crate::local::discovery;
use crate::local::docker;
use serde::{Deserialize, Serialize};
use std::path::PathBuf;
const DEFAULT_HTTP_PORT: u16 = 8123;
const DEFAULT_TCP_PORT: u16 = 9000;
const ADJECTIVES: &[&str] = &[
"bold", "calm", "dark", "fast", "gold", "keen", "loud", "neat", "pale", "red", "slim", "tall",
"warm", "blue", "cool", "deep", "flat", "gray", "iron", "wild",
];
const NOUNS: &[&str] = &[
"bear", "bird", "bolt", "crab", "crow", "dart", "fawn", "fish", "frog", "gull", "hare", "hawk",
"lynx", "moth", "newt", "orca", "puma", "seal", "swan", "wolf",
];
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Engine {
Clickhouse,
Postgres,
}
impl Engine {
pub fn as_str(&self) -> &'static str {
match self {
Engine::Clickhouse => "clickhouse",
Engine::Postgres => "postgres",
}
}
}
fn default_engine() -> Engine {
Engine::Clickhouse
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ServerInfo {
pub name: String,
pub pid: u32,
pub version: String,
pub http_port: u16,
pub tcp_port: u16,
pub started_at: String,
pub cwd: String,
#[serde(default = "default_engine")]
pub engine: Engine,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub container_id: Option<String>,
}
pub struct ServerEntry {
pub name: String,
pub running: bool,
pub info: Option<ServerInfo>,
}
pub fn validate_server_name(name: &str) -> Result<()> {
if name.is_empty()
|| name.contains('/')
|| name.contains('\\')
|| name.contains('\0')
|| name == "."
|| name == ".."
|| name.contains("../")
|| name.contains("..\\")
{
return Err(Error::InvalidServerName(name.to_string()));
}
Ok(())
}
fn servers_dir() -> PathBuf {
init::local_dir().join("servers")
}
fn server_meta_path(name: &str) -> PathBuf {
servers_dir().join(format!("{}.json", name))
}
pub fn server_meta_path_for_recovery(name: &str) -> PathBuf {
server_meta_path(name)
}
pub fn server_data_dir(name: &str) -> PathBuf {
servers_dir().join(name).join("data")
}
pub fn pg_instance_key(name: &str, major: &str) -> String {
format!("{}-pg{}", name, major)
}
pub fn servers_dir_join(child: &str) -> PathBuf {
servers_dir().join(child)
}
pub fn pg_data_dir(name: &str, major: &str) -> PathBuf {
servers_dir().join(pg_instance_key(name, major)).join("data")
}
fn ensure_servers_dir() -> Result<()> {
let dir = servers_dir();
if !dir.exists() {
std::fs::create_dir_all(&dir)?;
let gitignore = init::local_dir().join(".gitignore");
if !gitignore.exists() {
let _ = std::fs::write(gitignore, "*\n");
}
}
Ok(())
}
pub fn ensure_server_data_dir(name: &str) -> Result<()> {
ensure_servers_dir()?;
std::fs::create_dir_all(server_data_dir(name))?;
Ok(())
}
pub fn ensure_pg_data_dir(name: &str, major: &str) -> Result<()> {
ensure_servers_dir()?;
std::fs::create_dir_all(pg_data_dir(name, major))?;
Ok(())
}
pub fn save_server_info(info: &ServerInfo) -> Result<()> {
let dir = servers_dir();
std::fs::create_dir_all(&dir)?;
let path = server_meta_path(&info.name);
let json = serde_json::to_string_pretty(info)?;
std::fs::write(path, json)?;
Ok(())
}
pub fn remove_server_info(name: &str) {
let _ = std::fs::remove_file(server_meta_path(name));
}
fn is_alive(info: &ServerInfo) -> bool {
match info.engine {
Engine::Clickhouse => is_process_alive(info.pid),
Engine::Postgres => match info.container_id.as_deref() {
Some(id) => docker::is_container_running_blocking(id),
None => false,
},
}
}
pub fn load_info(name: &str) -> Option<ServerInfo> {
let content = std::fs::read_to_string(server_meta_path(name)).ok()?;
serde_json::from_str(&content).ok()
}
pub fn find_pg_instances(name: &str) -> Vec<ServerInfo> {
let prefix = format!("{}-pg", name);
let dir = match std::fs::read_dir(servers_dir()) {
Ok(d) => d,
Err(_) => return Vec::new(),
};
let mut out = Vec::new();
for entry in dir.flatten() {
let fname = match entry.file_name().into_string() {
Ok(s) => s,
Err(_) => continue,
};
let stem = match fname.strip_suffix(".json") {
Some(s) => s,
None => continue,
};
if !stem.starts_with(&prefix) {
continue;
}
let major = &stem[prefix.len()..];
if major.is_empty() || !major.chars().all(|c| c.is_ascii_digit()) {
continue;
}
if let Some(info) = load_info(stem)
&& info.engine == Engine::Postgres
{
out.push(info);
}
}
out
}
fn load_running_info(name: &str) -> Option<ServerInfo> {
let info = load_info(name)?;
if is_alive(&info) { Some(info) } else { None }
}
pub fn list_all_servers() -> Vec<ServerEntry> {
recover_current_project_servers();
let dir = servers_dir();
let mut entries = Vec::new();
let dir_entries = match std::fs::read_dir(&dir) {
Ok(e) => e,
Err(_) => return entries,
};
for entry in dir_entries.flatten() {
let path = entry.path();
let fname = match entry.file_name().into_string() {
Ok(s) => s,
Err(_) => continue,
};
let stem = match fname.strip_suffix(".json") {
Some(s) => s,
None => continue,
};
if !path.is_file() {
continue;
}
let info = load_info(stem);
let running = match &info {
Some(i) => is_alive(i),
None => false,
};
if let Some(i) = &info
&& !running
&& i.engine == Engine::Clickhouse
{
let _ = std::fs::remove_file(server_meta_path(stem));
entries.push(ServerEntry {
name: stem.to_string(),
running: false,
info: None,
});
continue;
}
entries.push(ServerEntry {
name: stem.to_string(),
running,
info,
});
}
entries.sort_by(|a, b| b.running.cmp(&a.running).then(a.name.cmp(&b.name)));
entries
}
pub fn list_running_servers() -> Vec<ServerInfo> {
list_all_servers()
.into_iter()
.filter(|e| e.running)
.filter_map(|e| e.info)
.collect()
}
pub fn is_server_running(name: &str) -> bool {
load_running_info(name).is_some()
}
pub fn running_server_count() -> usize {
list_running_servers().len()
}
fn is_process_alive(pid: u32) -> bool {
unsafe { libc::kill(pid as i32, 0) == 0 }
}
fn send_signal(pid: u32, signal: i32) -> Result<()> {
let ret = unsafe { libc::kill(pid as i32, signal) };
if ret != 0 {
let err = std::io::Error::last_os_error();
Err(Error::Exec(format!(
"Failed to send signal to PID {}: {}",
pid, err
)))
} else {
Ok(())
}
}
fn kill_process(pid: u32) -> Result<()> {
send_signal(pid, libc::SIGTERM)?;
std::thread::sleep(std::time::Duration::from_millis(500));
if is_process_alive(pid) {
std::thread::sleep(std::time::Duration::from_secs(2));
if is_process_alive(pid) {
send_signal(pid, libc::SIGKILL)?;
std::thread::sleep(std::time::Duration::from_millis(100));
}
}
if is_process_alive(pid) {
return Err(Error::Exec(format!(
"Process {} did not exit after SIGKILL",
pid
)));
}
Ok(())
}
pub fn kill_server(name: &str) -> Result<()> {
let info = load_running_info(name).ok_or_else(|| Error::ServerNotRunning(name.to_string()))?;
match info.engine {
Engine::Clickhouse => {
kill_process(info.pid)?;
remove_server_info(name);
}
Engine::Postgres => {
let id = info.container_id.as_deref().ok_or_else(|| {
Error::DockerError(format!(
"Postgres server '{}' has no container_id in metadata",
name
))
})?;
docker::stop_blocking(id)?;
}
}
Ok(())
}
pub fn resolve_name(name: Option<&str>) -> Result<String> {
match name {
Some(n) => {
validate_server_name(n)?;
Ok(n.to_string())
}
None => {
if is_server_running("default") {
Ok(generate_random_name())
} else {
Ok("default".to_string())
}
}
}
}
fn generate_random_name() -> String {
let seed = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
let mixed = seed ^ (std::process::id() as u128);
let adj = ADJECTIVES[(mixed % ADJECTIVES.len() as u128) as usize];
let noun = NOUNS[((mixed / ADJECTIVES.len() as u128) % NOUNS.len() as u128) as usize];
let tag = format!("{}-{}", adj, noun);
if is_server_running(&tag) {
for i in 2..100 {
let candidate = format!("{}-{}", tag, i);
if !is_server_running(&candidate) {
return candidate;
}
}
}
tag
}
pub fn check_spawn_health(pid: u32, name: &str) -> Result<()> {
std::thread::sleep(std::time::Duration::from_millis(300));
if !is_process_alive(pid) {
remove_server_info(name);
return Err(Error::Exec(format!(
"Server '{}' exited immediately after starting. \
Check if another server is using the same ports, \
or run in foreground to see the error output.",
name
)));
}
Ok(())
}
fn is_port_available(port: u16) -> bool {
std::net::TcpListener::bind(("127.0.0.1", port)).is_ok()
}
fn find_free_port(start: u16) -> Option<u16> {
(start..=start.saturating_add(100)).find(|&p| is_port_available(p))
}
pub fn resolve_ports(http_port: Option<u16>, tcp_port: Option<u16>) -> Result<(u16, u16, bool)> {
let http = match http_port {
Some(p) => p,
None => {
if is_port_available(DEFAULT_HTTP_PORT) {
DEFAULT_HTTP_PORT
} else {
find_free_port(DEFAULT_HTTP_PORT + 1)
.ok_or_else(|| Error::Exec("Could not find a free HTTP port".into()))?
}
}
};
let tcp = match tcp_port {
Some(p) => p,
None => {
if is_port_available(DEFAULT_TCP_PORT) {
DEFAULT_TCP_PORT
} else {
find_free_port(DEFAULT_TCP_PORT + 1)
.ok_or_else(|| Error::Exec("Could not find a free TCP port".into()))?
}
}
};
let auto_assigned = http_port.is_none() && http != DEFAULT_HTTP_PORT
|| tcp_port.is_none() && tcp != DEFAULT_TCP_PORT;
Ok((http, tcp, auto_assigned))
}
pub fn port_flags(http_port: u16, tcp_port: u16) -> Vec<String> {
vec![
format!("--http_port={}", http_port),
format!("--tcp_port={}", tcp_port),
]
}
pub fn now_timestamp() -> String {
let duration = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default();
format!("{}", duration.as_secs())
}
pub fn recover_current_project_servers() {
let current_dir = match std::env::current_dir().and_then(|p| p.canonicalize()) {
Ok(p) => p.display().to_string(),
Err(_) => return,
};
let processes = discovery::discover_clickhouse_processes();
for proc in processes {
let discovered_root = match std::path::Path::new(&proc.project_root).canonicalize() {
Ok(p) => p.display().to_string(),
Err(_) => proc.project_root.clone(),
};
if discovered_root != current_dir {
continue;
}
if load_running_info(&proc.server_name).is_some() {
continue;
}
let info = ServerInfo {
name: proc.server_name,
pid: proc.pid,
version: proc.version.unwrap_or_else(|| "unknown".to_string()),
http_port: proc.http_port.unwrap_or(0),
tcp_port: proc.tcp_port.unwrap_or(0),
started_at: "recovered".to_string(),
cwd: current_dir.clone(),
engine: Engine::Clickhouse,
container_id: None,
};
let _ = save_server_info(&info);
}
docker::recover_project_postgres_blocking(¤t_dir);
}
pub struct GlobalServerEntry {
pub name: String,
pub pid: u32,
pub project: String,
pub http_port: Option<u16>,
pub tcp_port: Option<u16>,
pub version: Option<String>,
pub engine: Engine,
pub container_id: Option<String>,
}
pub fn list_all_servers_global() -> Vec<GlobalServerEntry> {
let processes = discovery::discover_clickhouse_processes();
processes
.into_iter()
.map(|p| GlobalServerEntry {
name: p.server_name,
pid: p.pid,
project: p.project_root,
http_port: p.http_port,
tcp_port: p.tcp_port,
version: p.version,
engine: Engine::Clickhouse,
container_id: None,
})
.collect()
}
pub fn kill_server_by_pid(pid: u32) -> Result<()> {
if !is_process_alive(pid) {
return Err(Error::ServerNotRunning(format!("PID {}", pid)));
}
kill_process(pid)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn engine_serializes_lowercase() {
assert_eq!(serde_json::to_string(&Engine::Clickhouse).unwrap(), "\"clickhouse\"");
assert_eq!(serde_json::to_string(&Engine::Postgres).unwrap(), "\"postgres\"");
}
#[test]
fn server_info_legacy_json_deserializes_as_clickhouse() {
let legacy = r#"{
"name": "default",
"pid": 12345,
"version": "25.12.5.44",
"http_port": 8123,
"tcp_port": 9000,
"started_at": "1700000000",
"cwd": "/tmp/proj"
}"#;
let info: ServerInfo = serde_json::from_str(legacy).expect("legacy JSON should parse");
assert_eq!(info.engine, Engine::Clickhouse);
assert!(info.container_id.is_none());
}
#[test]
fn server_info_postgres_round_trip() {
let info = ServerInfo {
name: "dev".into(),
pid: 0,
version: "postgres:16".into(),
http_port: 0,
tcp_port: 5432,
started_at: "1700000000".into(),
cwd: "/tmp/proj".into(),
engine: Engine::Postgres,
container_id: Some("abc123".into()),
};
let json = serde_json::to_string(&info).unwrap();
let parsed: ServerInfo = serde_json::from_str(&json).unwrap();
assert_eq!(parsed.engine, Engine::Postgres);
assert_eq!(parsed.container_id.as_deref(), Some("abc123"));
assert!(json.contains("\"engine\":\"postgres\""));
}
}