use std::path::PathBuf;
use std::process::{Child, Command, Stdio};
use std::time::{Duration, Instant};
pub struct SupervisorConfig {
pub redis_server_bin: PathBuf,
pub falkordb_module: PathBuf,
pub data_dir: PathBuf,
pub port: u16,
pub max_restarts: u32,
pub port_file: Option<PathBuf>,
pub auth_token: Option<String>,
}
pub struct SupervisedServer {
pub child: Child,
pub port: u16,
pub restarts: u32,
auth_token: Option<String>,
}
impl Drop for SupervisedServer {
fn drop(&mut self) {
request_shutdown(self.port, self.auth_token.as_deref());
let deadline = Instant::now() + Duration::from_secs(2);
while Instant::now() < deadline {
if self.child.try_wait().ok().flatten().is_some() {
return;
}
std::thread::sleep(Duration::from_millis(20));
}
let _ = self.child.kill();
let _ = self.child.wait();
}
}
impl SupervisedServer {
pub fn poll(&mut self, cfg: &SupervisorConfig) -> anyhow::Result<()> {
match self.child.try_wait() {
Ok(Some(_status)) => {
if self.restarts >= cfg.max_restarts {
anyhow::bail!(
"supervised server exited; restart budget ({}) exhausted",
cfg.max_restarts
);
}
self.restarts += 1;
metrics::counter!("exocortex_supervisor_restarts_total").increment(1);
tracing::warn!(
restart = self.restarts,
port = self.port,
"supervised server crashed; restarting"
);
self.child = spawn_child(cfg)?;
if !wait_ping(cfg, &mut self.child)? {
anyhow::bail!("supervised server restart did not answer PING");
}
}
Ok(None) => {}
Err(e) => anyhow::bail!("supervisor try_wait failed: {e}"),
}
Ok(())
}
}
fn spawn_child(cfg: &SupervisorConfig) -> anyhow::Result<Child> {
let mut command = Command::new(&cfg.redis_server_bin);
command
.args([
"--port",
&cfg.port.to_string(),
"--bind",
"127.0.0.1",
"--save",
"1 1",
"--appendonly",
"yes",
"--appendfsync",
"everysec",
"--dir",
])
.arg(&cfg.data_dir)
.arg("--loadmodule")
.arg(&cfg.falkordb_module);
if let Some(token) = &cfg.auth_token {
command.arg("--requirepass").arg(token);
}
command
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.map_err(Into::into)
}
fn wait_ping(cfg: &SupervisorConfig, child: &mut Child) -> anyhow::Result<bool> {
let deadline = Instant::now() + Duration::from_secs(10);
loop {
if child.try_wait()?.is_some() {
anyhow::bail!("supervised redis-server exited during startup");
}
if ping(cfg.port, cfg.auth_token.as_deref()) {
return Ok(true);
}
if Instant::now() > deadline {
return Ok(false);
}
std::thread::sleep(Duration::from_millis(100));
}
}
pub fn resolve_paths(
flag_bin: Option<PathBuf>,
flag_module: Option<PathBuf>,
) -> anyhow::Result<(PathBuf, PathBuf)> {
let bin = flag_bin
.or_else(|| {
std::env::var("EXOCORTEX_REDIS_SERVER")
.ok()
.map(PathBuf::from)
})
.ok_or_else(|| {
anyhow::anyhow!(
"mcp-standalone needs a redis-server binary: pass --redis-server-bin \
or set EXOCORTEX_REDIS_SERVER (see crates/exocortex-server/src/supervisor.rs)"
)
})?;
let module = flag_module
.or_else(|| {
std::env::var("EXOCORTEX_FALKORDB_MODULE")
.ok()
.map(PathBuf::from)
})
.ok_or_else(|| {
anyhow::anyhow!(
"mcp-standalone needs the FalkorDB module path: pass --falkordb-module \
or set EXOCORTEX_FALKORDB_MODULE"
)
})?;
Ok((bin, module))
}
pub fn spawn_supervised(cfg: &SupervisorConfig) -> anyhow::Result<SupervisedServer> {
std::fs::create_dir_all(&cfg.data_dir)?;
let mut child = spawn_child(cfg)?;
if !wait_ping(cfg, &mut child)? {
let _ = child.kill();
anyhow::bail!("supervised FalkorDB server did not answer PING within 10s");
}
if let Some(path) = &cfg.port_file {
exocortex_storage::bounded_io::atomic_write_private(
path,
cfg.port.to_string().as_bytes(),
"supervised port",
)?;
}
tracing::info!(port = cfg.port, "supervised FalkorDB server up");
Ok(SupervisedServer {
child,
port: cfg.port,
restarts: 0,
auth_token: cfg.auth_token.clone(),
})
}
pub fn supervised_store_urls(port: u16, auth_token: Option<&str>) -> (String, String) {
let authority = auth_token
.map(|token| format!(":{token}@"))
.unwrap_or_default();
(
format!("falkor://{authority}127.0.0.1:{port}"),
format!("redis://{authority}127.0.0.1:{port}"),
)
}
fn ping(port: u16, auth_token: Option<&str>) -> bool {
use std::io::{Read, Write};
let Ok(mut s) = std::net::TcpStream::connect(("127.0.0.1", port)) else {
return false;
};
let mut buf = [0u8; 128];
if let Some(token) = auth_token {
if s.write_all(format!("AUTH {token}\r\n").as_bytes()).is_err() {
return false;
}
let Ok(n) = s.read(&mut buf) else {
return false;
};
if !buf[..n].starts_with(b"+OK") {
return false;
}
}
if s.write_all(b"PING\r\n").is_err() {
return false;
}
let Ok(n) = s.read(&mut buf) else {
return false;
};
buf[..n].windows(4).any(|w| w == b"PONG" || w == b"+PON")
}
fn request_shutdown(port: u16, auth_token: Option<&str>) {
use std::io::{Read as _, Write as _};
if let Ok(mut stream) = std::net::TcpStream::connect(("127.0.0.1", port)) {
if let Some(token) = auth_token {
let _ = stream.write_all(format!("AUTH {token}\r\n").as_bytes());
let _ = stream.read(&mut [0u8; 64]);
}
let _ = stream.write_all(b"SHUTDOWN SAVE\r\n");
}
}
pub fn free_port() -> anyhow::Result<u16> {
let listener = std::net::TcpListener::bind(("127.0.0.1", 0))?;
Ok(listener.local_addr()?.port())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn supervise_restarts_within_budget_then_gives_up() {
let cfg = SupervisorConfig {
redis_server_bin: "/bin/sleep".into(),
falkordb_module: "unused".into(),
data_dir: std::env::temp_dir(),
port: 0,
max_restarts: 2,
port_file: None,
auth_token: None,
};
let mut server = SupervisedServer {
child: Command::new("/bin/sleep")
.arg("1")
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap(),
port: 0,
restarts: 0,
auth_token: None,
};
server.child.kill().unwrap();
let _ = server.child.wait();
let mut restarts = 0;
loop {
let crashed = matches!(server.child.try_wait(), Ok(Some(_)));
if !crashed {
break;
}
if restarts >= cfg.max_restarts {
break;
}
restarts += 1;
server.child = Command::new("/bin/sleep")
.arg("0") .stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap();
let _ = server.child.wait();
}
assert_eq!(restarts, cfg.max_restarts, "restart policy bounds the loop");
}
#[test]
fn supervised_server_kills_child_on_drop() {
let child = Command::new("/bin/sleep")
.arg("30")
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap();
let pid = child.id();
let server = SupervisedServer {
child,
port: 0,
restarts: 0,
auth_token: None,
};
drop(server);
std::thread::sleep(Duration::from_millis(100));
let gone = Command::new("kill")
.arg("-0")
.arg(pid.to_string())
.stdout(Stdio::null())
.stderr(Stdio::null())
.status()
.map(|s| !s.success())
.unwrap_or(true);
assert!(gone, "drop killed the supervised child");
}
#[test]
fn free_port_returns_open_port() {
let port = free_port().unwrap();
assert!(port > 0);
}
#[test]
fn supervised_store_urls_embed_the_per_boot_token() {
let (falkor, redis) = supervised_store_urls(16379, Some("a1b2c3"));
assert_eq!(falkor, "falkor://:a1b2c3@127.0.0.1:16379");
assert_eq!(redis, "redis://:a1b2c3@127.0.0.1:16379");
let (falkor, redis) = supervised_store_urls(16379, None);
assert_eq!(falkor, "falkor://127.0.0.1:16379");
assert_eq!(redis, "redis://127.0.0.1:16379");
}
#[test]
fn supervised_store_rejects_unauthenticated_local_peers() {
let (Ok(bin), Ok(module)) = (
std::env::var("EXOCORTEX_REDIS_SERVER"),
std::env::var("EXOCORTEX_FALKORDB_MODULE"),
) else {
eprintln!(
"SKIP supervised_store_rejects_unauthenticated_local_peers: \
EXOCORTEX_REDIS_SERVER/EXOCORTEX_FALKORDB_MODULE absent; live suite unexecuted"
);
return;
};
let data_dir =
std::env::temp_dir().join(format!("exocortex-supervisor-auth-{}", std::process::id()));
let cfg = SupervisorConfig {
redis_server_bin: bin.into(),
falkordb_module: module.into(),
data_dir: data_dir.clone(),
port: free_port().unwrap(),
max_restarts: 0,
port_file: None,
auth_token: Some("5f4d3c2b1a5f4d3c2b1a5f4d3c2b1a5f4d3c2b1a5f4d3c2b1a".into()),
};
let server = spawn_supervised(&cfg).expect("supervised server with auth starts");
let refused = {
use std::io::{Read, Write};
let mut stream = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
stream.write_all(b"PING\r\n").unwrap();
let mut reply = [0u8; 64];
let n = stream.read(&mut reply).unwrap();
reply[..n].windows(6).any(|window| window == b"NOAUTH")
};
assert!(refused, "an unauthenticated local peer is refused");
assert!(
ping(server.port, cfg.auth_token.as_deref()),
"the authenticated handshake still answers"
);
drop(server);
let _ = std::fs::remove_dir_all(data_dir);
}
}