use std::collections::BTreeMap;
use std::fs;
use std::net::{SocketAddr, TcpStream};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread;
use std::time::{Duration, SystemTime};
use crate::project::ProjectConfig;
use super::{CancellableOutcome, ChildGuard, ProcessError, ProcessSpec, run_cancellable, spawn};
pub(crate) fn supervise(project: &ProjectConfig) -> Result<(), ProcessError> {
let stopping = install_shutdown_listener()?;
if matches!(
build_backend(project, &stopping)?,
CancellableOutcome::Cancelled
) {
return Ok(());
}
let mut frontend = start_frontend(project)?;
let mut backend = Some(start_backend(project)?);
wait_ready("frontend", project.frontend_port, &mut frontend, &stopping)?;
if let Some(child) = backend.as_mut() {
wait_ready("backend", project.backend_port, child, &stopping)?;
}
println!("backend ready http://127.0.0.1:{}", project.backend_port);
println!(
"frontend ready http://127.0.0.1:{}",
project.frontend_port
);
let mut snapshot = rust_snapshot(project.root())?;
let result = loop {
if stopping.load(Ordering::SeqCst) {
break Ok(());
}
if let Some(status) = frontend.try_wait()? {
if stopping.load(Ordering::SeqCst) {
break Ok(());
}
break Err(ProcessError::Failed {
program: "node (Vite)".to_owned(),
directory: project.frontend_root(),
status,
});
}
if let Some(child) = backend.as_mut()
&& let Some(status) = child.try_wait()?
{
if stopping.load(Ordering::SeqCst) {
break Ok(());
}
break Err(ProcessError::Failed {
program: project.backend_binary.clone(),
directory: project.root().to_path_buf(),
status,
});
}
thread::sleep(Duration::from_millis(250));
let next = rust_snapshot(project.root())?;
if next != snapshot {
snapshot = next;
if let Some(child) = backend.as_mut() {
child.terminate()?;
}
backend = None;
println!("backend rebuilding");
match build_backend(project, &stopping) {
Ok(CancellableOutcome::Completed) => {}
Ok(CancellableOutcome::Cancelled) => break Ok(()),
Err(error) => {
eprintln!("backend build failed: {error}");
continue;
}
}
let mut child = start_backend(project)?;
wait_ready("backend", project.backend_port, &mut child, &stopping)?;
backend = Some(child);
println!("backend ready http://127.0.0.1:{}", project.backend_port);
}
};
let backend_stop = match backend.as_mut() {
Some(child) => child.terminate(),
None => Ok(()),
};
let frontend_stop = frontend.terminate();
result.and(backend_stop).and(frontend_stop)
}
fn install_shutdown_listener() -> Result<Arc<AtomicBool>, ProcessError> {
let stopping = Arc::new(AtomicBool::new(false));
let signal = Arc::clone(&stopping);
let (sender, receiver) = std::sync::mpsc::sync_channel(1);
thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_io()
.build();
match runtime {
Ok(runtime) => {
if sender.send(Ok(())).is_ok() && runtime.block_on(wait_for_shutdown()).is_ok() {
signal.store(true, Ordering::SeqCst);
}
}
Err(error) => {
let _ = sender.send(Err(error));
}
}
});
receiver
.recv()
.map_err(|error| ProcessError::Signal(std::io::Error::other(error)))?
.map_err(ProcessError::Signal)?;
Ok(stopping)
}
#[cfg(unix)]
async fn wait_for_shutdown() -> std::io::Result<()> {
let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())?;
tokio::select! { result = tokio::signal::ctrl_c() => result, _ = terminate.recv() => Ok(()) }
}
#[cfg(windows)]
async fn wait_for_shutdown() -> std::io::Result<()> {
let mut control_break = tokio::signal::windows::ctrl_break()?;
tokio::select! {
result = tokio::signal::ctrl_c() => result,
_ = control_break.recv() => Ok(()),
}
}
#[cfg(not(any(unix, windows)))]
async fn wait_for_shutdown() -> std::io::Result<()> {
tokio::signal::ctrl_c().await
}
fn build_backend(
project: &ProjectConfig,
stopping: &AtomicBool,
) -> Result<CancellableOutcome, ProcessError> {
run_cancellable(
&ProcessSpec::new("cargo", project.root())
.args(["build", "--package", &project.backend_package])
.new_process_group(true),
stopping,
)
}
fn start_backend(project: &ProjectConfig) -> Result<ChildGuard, ProcessError> {
let extension = if cfg!(windows) { ".exe" } else { "" };
let binary = project
.root()
.join("target")
.join("debug")
.join(format!("{}{extension}", project.backend_binary));
spawn(
&ProcessSpec::new(binary.into_os_string(), project.root())
.env(
"ARCATURE_VITE_ORIGIN",
format!("http://127.0.0.1:{}", project.frontend_port),
)
.env("ARCATURE_BACKEND_PORT", project.backend_port.to_string())
.new_process_group(true),
)
.map(|child| ChildGuard::new(child, "backend"))
}
fn start_frontend(project: &ProjectConfig) -> Result<ChildGuard, ProcessError> {
let vite = project
.frontend_root()
.join("node_modules")
.join("vite")
.join("bin")
.join("vite.js");
let port = project.frontend_port.to_string();
spawn(
&ProcessSpec::new("node", project.frontend_root())
.arg(vite.into_os_string())
.args(["--host", "127.0.0.1", "--port", &port, "--strictPort"])
.env("ARCATURE_FRONTEND_PORT", &port)
.new_process_group(true),
)
.map(|child| ChildGuard::new(child, "frontend"))
}
fn wait_ready(
service: &'static str,
port: u16,
child: &mut ChildGuard,
stopping: &AtomicBool,
) -> Result<(), ProcessError> {
let address = SocketAddr::from(([127, 0, 0, 1], port));
for _ in 0..100 {
if stopping.load(Ordering::SeqCst) {
return Ok(());
}
if TcpStream::connect_timeout(&address, Duration::from_millis(100)).is_ok() {
return Ok(());
}
if let Some(status) = child.try_wait()? {
if stopping.load(Ordering::SeqCst) {
return Ok(());
}
return Err(ProcessError::Failed {
program: service.to_owned(),
directory: PathBuf::from("."),
status,
});
}
thread::sleep(Duration::from_millis(100));
}
Err(ProcessError::Readiness {
service,
address: address.to_string(),
})
}
fn rust_snapshot(root: &Path) -> Result<BTreeMap<PathBuf, (SystemTime, u64)>, ProcessError> {
let mut snapshot = BTreeMap::new();
collect(&root.join("src"), &mut snapshot)?;
for name in ["Cargo.toml", "Cargo.lock", "build.rs"] {
let path = root.join(name);
if path.is_file() {
insert_metadata(&path, &mut snapshot)?;
}
}
Ok(snapshot)
}
fn collect(
path: &Path,
snapshot: &mut BTreeMap<PathBuf, (SystemTime, u64)>,
) -> Result<(), ProcessError> {
let entries = match fs::read_dir(path) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(source) => {
return Err(ProcessError::Inspect {
path: path.to_path_buf(),
source,
});
}
};
for entry in entries {
let entry = entry.map_err(|source| ProcessError::Inspect {
path: path.to_path_buf(),
source,
})?;
let child = entry.path();
let kind = entry.file_type().map_err(|source| ProcessError::Inspect {
path: child.clone(),
source,
})?;
if kind.is_dir() {
collect(&child, snapshot)?;
} else if kind.is_file() && child.extension().is_some_and(|extension| extension == "rs") {
insert_metadata(&child, snapshot)?;
}
}
Ok(())
}
fn insert_metadata(
path: &Path,
snapshot: &mut BTreeMap<PathBuf, (SystemTime, u64)>,
) -> Result<(), ProcessError> {
let metadata = fs::metadata(path).map_err(|source| ProcessError::Inspect {
path: path.to_path_buf(),
source,
})?;
let modified = metadata
.modified()
.map_err(|source| ProcessError::Inspect {
path: path.to_path_buf(),
source,
})?;
snapshot.insert(path.to_path_buf(), (modified, metadata.len()));
Ok(())
}