use crate::libs::config::Config;
use crate::libs::data_storage::DataStorage;
use crate::libs::messages::Message;
use crate::libs::monitor::Monitor;
use crate::{msg_bail_anyhow, msg_error, msg_error_anyhow, msg_info, msg_warning};
use anyhow::Result;
use std::time::Duration;
use tracing::{debug, info, instrument, warn};
const PID_FILE: &str = "kasl-watch.pid";
#[instrument]
pub async fn run_with_signal_handling() -> Result<()> {
info!("Starting daemon with signal handling");
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
#[cfg(unix)]
{
tokio::spawn(async move {
use tokio::signal::unix::{SignalKind, signal};
let mut sigterm = signal(SignalKind::terminate()).unwrap_or_else(|_| panic!("{}", Message::FailedToCreateSigtermHandler));
let mut sigint = signal(SignalKind::interrupt()).unwrap_or_else(|_| panic!("{}", Message::FailedToCreateSigintHandler));
tokio::select! {
_ = sigterm.recv() => {
msg_info!(Message::WatcherReceivedSigterm);
}
_ = sigint.recv() => {
msg_info!(Message::WatcherReceivedSigint);
}
}
let _ = shutdown_tx.send(());
});
}
#[cfg(windows)]
{
tokio::spawn(async move {
match tokio::signal::ctrl_c().await {
Ok(()) => {
msg_info!(Message::WatcherReceivedCtrlC);
}
Err(e) => {
msg_error!(Message::WatcherCtrlCListenFailed(e.to_string()));
}
}
let _ = shutdown_tx.send(());
});
}
#[cfg(not(any(unix, windows)))]
{
msg_warning!(Message::WatcherSignalHandlingNotSupported);
}
let monitor_handle = tokio::spawn(async move {
match run_monitor().await {
Ok(()) => Ok(()),
Err(e) => Err(Message::MonitorError(e.to_string())),
}
});
let inbox_handle = tokio::spawn(async move {
crate::libs::jira_inbox::run_poller().await;
});
tokio::select! {
result = monitor_handle => {
inbox_handle.abort();
match result {
Ok(Ok(())) => msg_info!(Message::MonitorExitedNormally),
Ok(Err(e)) => msg_error!(Message::MonitorError(e.to_string())),
Err(e) => msg_error!(Message::MonitorTaskPanicked(e.to_string())),
}
}
_ = shutdown_rx => {
inbox_handle.abort();
msg_info!(Message::MonitorShuttingDown);
}
}
let pid_path = DataStorage::new().get_path(PID_FILE)?;
if pid_path.exists() {
let _ = std::fs::remove_file(&pid_path);
}
Ok(())
}
async fn run_monitor() -> Result<()> {
let config = Config::read()?;
let monitor_config = config.monitor.unwrap_or_default();
let mut monitor = Monitor::new(monitor_config)?;
monitor.run().await
}
#[instrument]
pub fn spawn() -> Result<()> {
debug!("Attempting to spawn daemon process");
let pid_path = DataStorage::new().get_path(PID_FILE)?;
if pid_path.exists()
&& let Ok(pid_str) = std::fs::read_to_string(&pid_path)
{
msg_info!(Message::WatcherStoppingExisting(pid_str.trim().to_string()));
if let Err(e) = stop_internal() {
msg_warning!(Message::WatcherFailedToStopExisting(e.to_string()));
let _ = std::fs::remove_file(&pid_path);
}
std::thread::sleep(Duration::from_millis(1000));
}
let current_exe = std::env::current_exe().unwrap_or_else(|_| panic!("{}", Message::FailedToGetCurrentExecutable.to_string()));
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
let mut command = std::process::Command::new(current_exe);
command.arg("--daemon-run");
unsafe {
command.pre_exec(|| {
nix::unistd::setsid()?;
Ok(())
});
}
let child = command.spawn()?;
let pid = child.id();
std::fs::write(pid_path, pid.to_string())?;
msg_info!(Message::WatcherStarted(pid));
}
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
const CREATE_NO_WINDOW: u32 = 0x08000000;
let child = std::process::Command::new(current_exe)
.arg("--daemon-run")
.creation_flags(CREATE_NO_WINDOW)
.spawn()?;
let pid = child.id();
std::fs::write(pid_path, pid.to_string())?;
msg_info!(Message::WatcherStarted(pid));
}
#[cfg(not(any(unix, windows)))]
{
msg_bail_anyhow!(Message::DaemonModeNotSupported);
}
Ok(())
}
pub fn is_running() -> bool {
let pid_path = match DataStorage::new().get_path(PID_FILE) {
Ok(path) => path,
Err(_) => return false,
};
if !pid_path.exists() {
return false;
}
let pid_str = match std::fs::read_to_string(&pid_path) {
Ok(content) => content,
Err(_) => return false,
};
let pid: u32 = match pid_str.trim().parse() {
Ok(pid) => pid,
Err(_) => return false,
};
is_process_running(pid)
}
fn is_process_running(pid: u32) -> bool {
#[cfg(windows)]
{
use winapi::um::errhandlingapi::GetLastError;
use winapi::um::handleapi::CloseHandle;
use winapi::um::processthreadsapi::OpenProcess;
use winapi::um::winnt::PROCESS_QUERY_INFORMATION;
unsafe {
let handle = OpenProcess(PROCESS_QUERY_INFORMATION, 0, pid);
if handle.is_null() {
let error = GetLastError();
return error != 87;
}
CloseHandle(handle);
true
}
}
#[cfg(unix)]
{
use std::process::Command;
match Command::new("ps").arg("-p").arg(pid.to_string()).output() {
Ok(output) => output.status.success(),
Err(_) => false,
}
}
#[cfg(not(any(unix, windows)))]
{
false
}
}
pub fn stop() -> Result<()> {
match stop_internal() {
Ok(()) => Ok(()),
Err(e) => {
if e.to_string().contains("not found") || e.to_string().contains("not running") {
msg_info!(Message::WatcherNotRunning);
Ok(())
} else {
Err(e)
}
}
}
}
fn stop_internal() -> Result<()> {
let pid_path = DataStorage::new().get_path(PID_FILE)?;
let pid_str = match std::fs::read_to_string(&pid_path) {
Ok(content) => content,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
msg_bail_anyhow!(Message::WatcherNotRunningPidNotFound);
}
Err(e) => return Err(e.into()),
};
let pid: u32 = pid_str.trim().parse().map_err(|_| msg_error_anyhow!(Message::InvalidPidFileContent))?;
let killed = kill_process(pid)?;
if let Err(e) = std::fs::remove_file(&pid_path)
&& e.kind() != std::io::ErrorKind::NotFound
{
return Err(e.into());
}
if killed {
msg_info!(Message::WatcherStopped(pid));
} else {
msg_info!(Message::WatcherNotRunning);
}
Ok(())
}
#[cfg(windows)]
fn kill_process(pid: u32) -> Result<bool> {
use winapi::um::errhandlingapi::GetLastError;
use winapi::um::handleapi::CloseHandle;
use winapi::um::processthreadsapi::{OpenProcess, TerminateProcess};
use winapi::um::winnt::PROCESS_TERMINATE;
unsafe {
let handle = OpenProcess(PROCESS_TERMINATE, 0, pid);
if handle.is_null() {
let error = GetLastError();
if error == 87 {
return Ok(false);
}
msg_bail_anyhow!(Message::FailedToOpenProcess(error));
}
let result = TerminateProcess(handle, 0);
CloseHandle(handle);
if result == 0 {
let error = GetLastError();
msg_bail_anyhow!(Message::FailedToTerminateProcess(error));
} else {
std::thread::sleep(Duration::from_millis(100));
Ok(true)
}
}
}
#[cfg(unix)]
fn kill_process(pid: u32) -> Result<bool> {
use std::process::Command;
let output = Command::new("ps").arg("-p").arg(pid.to_string()).output()?;
if !output.status.success() {
return Ok(false);
}
Command::new("kill").arg("-TERM").arg(pid.to_string()).output()?;
for _ in 0..10 {
std::thread::sleep(Duration::from_millis(100));
let check = Command::new("ps").arg("-p").arg(pid.to_string()).output()?;
if !check.status.success() {
return Ok(true);
}
}
Command::new("kill").arg("-9").arg(pid.to_string()).output()?;
std::thread::sleep(Duration::from_millis(100));
Ok(true)
}
#[cfg(not(any(unix, windows)))]
fn kill_process(_pid: u32) -> Result<bool> {
msg_bail_anyhow!(Message::ProcessTerminationNotSupported);
}