use std::io::{Read, Write};
use std::net::{SocketAddr, TcpStream};
use std::path::{Path, PathBuf};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use super::model::DEFAULT_PORT;
pub const PORT_ENV: &str = "SCSH_DAEMON_PORT";
pub const HOME_ENV: &str = "SCSH_HOME";
pub fn scsh_home_dir() -> PathBuf {
if let Some(dir) = std::env::var_os(HOME_ENV).filter(|s| !s.is_empty()) {
return PathBuf::from(dir);
}
match std::env::var_os("HOME").filter(|s| !s.is_empty()) {
Some(home) => PathBuf::from(home).join(".scsh"),
None => daemon_dir(),
}
}
pub fn store_db_file(port: u16) -> PathBuf {
scsh_home_dir().join(format!("daemon-{port}.redb"))
}
pub fn now_unix_secs() -> u64 {
SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0)
}
pub fn daemon_port() -> u16 {
std::env::var(PORT_ENV).ok().and_then(|s| s.parse().ok()).unwrap_or(DEFAULT_PORT)
}
pub fn daemon_dir() -> PathBuf {
std::env::temp_dir().join("scsh-daemon")
}
pub fn pid_file(port: u16) -> PathBuf {
daemon_dir().join(format!("daemon-{port}.pid"))
}
pub fn prune_file(port: u16) -> PathBuf {
daemon_dir().join(format!("prune-{port}.json"))
}
pub fn mode_file(port: u16) -> PathBuf {
daemon_dir().join(format!("daemon-{port}.mode"))
}
pub fn proc_restart_marker(session_id: &str, proc_index: usize) -> PathBuf {
daemon_dir().join("restart-requests").join(format!("{session_id}-{proc_index}"))
}
pub fn request_proc_restart(session_id: &str, proc_index: usize) -> bool {
let path = proc_restart_marker(session_id, proc_index);
path.parent().is_some_and(|dir| std::fs::create_dir_all(dir).is_ok()) && std::fs::write(&path, b"").is_ok()
}
pub fn consume_proc_restart(session_id: &str, proc_index: usize) -> bool {
std::fs::remove_file(proc_restart_marker(session_id, proc_index)).is_ok()
}
pub fn session_cancel_marker(session_id: &str) -> PathBuf {
daemon_dir().join("cancel-requests").join(session_id)
}
pub fn request_session_cancel(session_id: &str) -> bool {
let path = session_cancel_marker(session_id);
path.parent().is_some_and(|dir| std::fs::create_dir_all(dir).is_ok()) && std::fs::write(&path, b"").is_ok()
}
pub fn session_cancelled(session_id: &str) -> bool {
!session_id.is_empty() && session_cancel_marker(session_id).exists()
}
pub fn clear_session_cancel(session_id: &str) {
let _ = std::fs::remove_file(session_cancel_marker(session_id));
}
pub fn daemon_port_reachable(port: u16) -> bool {
let addr: SocketAddr = format!("127.0.0.1:{port}").parse().expect("valid localhost address");
TcpStream::connect_timeout(&addr, Duration::from_millis(200)).is_ok()
}
pub fn daemon_api_responds(port: u16) -> bool {
if !daemon_port_reachable(port) {
return false;
}
let addr: SocketAddr = format!("127.0.0.1:{port}").parse().expect("valid localhost address");
let Ok(mut stream) = TcpStream::connect_timeout(&addr, Duration::from_millis(500)) else {
return false;
};
stream.set_read_timeout(Some(Duration::from_secs(2))).ok();
stream.set_write_timeout(Some(Duration::from_secs(2))).ok();
let req = "GET /api/v1/sessions HTTP/1.1\r\nHost: 127.0.0.1\r\nConnection: close\r\n\r\n";
if stream.write_all(req.as_bytes()).is_err() {
return false;
}
let mut resp = String::new();
if stream.read_to_string(&mut resp).is_err() {
return false;
}
resp.starts_with("HTTP/1.1 200") && resp.contains("application/json")
}
pub fn daemon_reported_version(port: u16) -> Option<String> {
let body = fetch_version_body(port)?;
let after = body.split("\"version\":").nth(1)?;
let start = after.find('"')? + 1;
let end = after[start..].find('"')? + start;
Some(after[start..end].to_string())
}
pub fn daemon_reported_started_at(port: u16) -> Option<u64> {
let body = fetch_version_body(port)?;
let after = body.split("\"started_at\":").nth(1)?;
let digits: String = after.trim_start().chars().take_while(|c| c.is_ascii_digit()).collect();
digits.parse().ok()
}
fn fetch_version_body(port: u16) -> Option<String> {
daemon_get_body(port, "/api/v1/version")
}
pub fn daemon_get_body(port: u16, path: &str) -> Option<String> {
let addr: SocketAddr = format!("127.0.0.1:{port}").parse().ok()?;
let mut stream = TcpStream::connect_timeout(&addr, Duration::from_millis(500)).ok()?;
stream.set_read_timeout(Some(Duration::from_secs(2))).ok();
stream.set_write_timeout(Some(Duration::from_secs(2))).ok();
let req = format!("GET {path} HTTP/1.1\r\nHost: 127.0.0.1\r\nConnection: close\r\n\r\n");
stream.write_all(req.as_bytes()).ok()?;
let mut resp = String::new();
stream.read_to_string(&mut resp).ok()?;
if !resp.starts_with("HTTP/1.1 200") {
return None;
}
resp.split("\r\n\r\n").nth(1).map(str::to_string)
}
pub fn read_persisted_mode(port: u16) -> Option<super::model::DaemonMode> {
let text = std::fs::read_to_string(mode_file(port)).ok()?;
super::model::DaemonMode::parse(text.trim())
}
pub fn write_mode_marker(port: u16, mode: super::model::DaemonMode) {
let _ = std::fs::create_dir_all(daemon_dir());
let _ = crate::atomic_write(&mode_file(port), mode.as_str().as_bytes());
}
pub fn signal_process(pid: u32, sig: i32) {
#[cfg(unix)]
{
unsafe {
libc::kill(pid as i32, sig);
}
}
}
pub fn pid_alive(pid: u32) -> bool {
if pid == 0 {
return false;
}
#[cfg(unix)]
{
unsafe { libc::kill(pid as i32, 0) == 0 }
}
#[cfg(not(unix))]
{
false
}
}
pub fn is_scsh_daemon_pid(pid: u32) -> bool {
#[cfg(unix)]
{
process_args(pid).is_some_and(|args| args.contains("scsh") && args.contains("__daemon-serve"))
}
#[cfg(not(unix))]
{
let _ = pid;
false
}
}
#[cfg(unix)]
fn process_args(pid: u32) -> Option<String> {
let output = std::process::Command::new("ps").arg("-p").arg(pid.to_string()).arg("-o").arg("args=").output().ok()?;
if !output.status.success() {
return None;
}
Some(String::from_utf8_lossy(&output.stdout).trim().to_string())
}
pub fn read_live_pid(port: u16) -> Option<u32> {
let text = std::fs::read_to_string(pid_file(port)).ok()?;
let pid = text.trim().parse::<u32>().ok()?;
if pid_alive(pid) && is_scsh_daemon_pid(pid) {
Some(pid)
} else {
None
}
}
pub fn projects_dir() -> std::path::PathBuf {
crate::runtime::scsh_home().join("projects")
}
pub fn base_url(port: u16) -> String {
format!("http://127.0.0.1:{port}")
}
pub fn session_url(port: u16, session_id: &str) -> String {
format!("{}/job/{}", base_url(port), session_id)
}
pub fn absolutize_repo_path(path: &Path) -> String {
let path = if path.is_absolute() {
path.to_path_buf()
} else {
std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")).join(path)
};
std::fs::canonicalize(&path).unwrap_or(path).to_string_lossy().into_owned()
}
#[cfg(unix)]
pub fn daemon_detach_child() -> std::io::Result<()> {
let pid = unsafe { libc::fork() };
if pid < 0 {
return Err(std::io::Error::last_os_error());
}
if pid > 0 {
unsafe { libc::_exit(0) };
}
let sid = unsafe { libc::setsid() };
if sid < 0 {
return Err(std::io::Error::last_os_error());
}
Ok(())
}
#[cfg(unix)]
mod libc {
#[link(name = "c")]
extern "C" {
pub fn kill(pid: i32, sig: i32) -> i32;
pub fn fork() -> i32;
pub fn setsid() -> i32;
pub fn _exit(code: i32) -> !;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn pid_file_names_include_port() {
let p = pid_file(7274);
assert!(p.to_string_lossy().contains("7274"));
assert!(p.to_string_lossy().ends_with(".pid"));
}
#[test]
fn session_url_format() {
let u = session_url(7274, "abcdef");
assert_eq!(u, "http://127.0.0.1:7274/job/abcdef");
}
#[test]
fn a_session_cancel_marker_is_peeked_by_every_route_and_cleared_once() {
let id = format!("cancel-test-{}", std::process::id());
clear_session_cancel(&id);
assert!(!session_cancelled(&id), "no marker, no cancellation");
assert!(request_session_cancel(&id));
assert!(session_cancelled(&id));
assert!(session_cancelled(&id));
assert!(session_cancelled(&id));
clear_session_cancel(&id);
assert!(!session_cancelled(&id), "teardown clears it");
clear_session_cancel(&id);
assert!(!session_cancelled(""));
}
}