use super::*;
pub async fn serve(repo_root: &Path) -> Result<()> {
let startup_t0 = std::time::Instant::now();
let mati_root: PathBuf = crate::store::mati_home_opt()
.map(|h| h.join(crate::store::derive_slug(repo_root)))
.ok_or_else(|| anyhow::anyhow!("cannot resolve home directory for mati_root"))?;
super::metadata::record_lifecycle_event(&mati_root, "startup", "phase=ensure_daemon");
if !super::daemon_lifecycle::ensure_daemon(&mati_root).await {
super::metadata::record_lifecycle_event(
&mati_root,
"serve_failed",
"daemon unreachable after auto-spawn",
);
anyhow::bail!(
"mati serve: daemon unreachable. \
Run `mati daemon start` manually and check the lifecycle.log."
);
}
super::metadata::record_lifecycle_event(
&mati_root,
"serve_start",
&format!("pid={} owner=proxy", std::process::id()),
);
super::metrics::init();
super::metadata::record_lifecycle_event(
&mati_root,
"startup",
&format!(
"phase=ready elapsed_ms={}",
startup_t0.elapsed().as_millis()
),
);
let worktree_tag = crate::store::session::worktree_scope_tag(repo_root);
let transport = rmcp::transport::io::stdio();
let service = MatiServer::with_socket_root(mati_root.clone(), worktree_tag)
.serve(transport)
.await
.map_err(|e| anyhow::anyhow!("MCP proxy initialization failed: {e}"))
.inspect_err(|e| {
super::metadata::record_lifecycle_event(
&mati_root,
"serve_failed",
&format!("proxy init: {e:#}"),
)
})?;
let shutdown_reason: &'static str = match service.waiting().await {
Ok(_) => "client_disconnect",
Err(e) => {
super::metadata::record_lifecycle_event(
&mati_root,
"serve_failed",
&format!("proxy waiting: {e}"),
);
"mcp_waiting_error"
}
};
super::metadata::record_lifecycle_event(
&mati_root,
"serve_shutdown",
&format!("reason={shutdown_reason}"),
);
Ok(())
}
pub(crate) async fn proxy_daemon_result(
root: &Path,
cmd: &str,
args: serde_json::Value,
) -> ProxyDaemonResult {
let result = proxy_daemon_result_no_spawn(root, cmd, &args).await;
if matches!(
&result,
ProxyDaemonResult::NotRunning | ProxyDaemonResult::StaleSocket
) && super::daemon_lifecycle::ensure_daemon(root).await
{
match proxy_daemon_result_once(root, cmd, &args).await {
AttemptOutcome::Final(r) | AttemptOutcome::Retryable(r) => return r,
}
}
result
}
pub(crate) async fn proxy_daemon_result_no_spawn(
root: &Path,
cmd: &str,
args: &serde_json::Value,
) -> ProxyDaemonResult {
match proxy_daemon_result_once(root, cmd, args).await {
AttemptOutcome::Final(result) => result,
AttemptOutcome::Retryable(_) => {
tokio::time::sleep(Duration::from_millis(100)).await;
match proxy_daemon_result_once(root, cmd, args).await {
AttemptOutcome::Final(result) | AttemptOutcome::Retryable(result) => result,
}
}
}
}
enum AttemptOutcome {
Final(ProxyDaemonResult),
Retryable(ProxyDaemonResult),
}
async fn proxy_daemon_result_once(
root: &Path,
cmd: &str,
args: &serde_json::Value,
) -> AttemptOutcome {
let v2_cmd = super::protocol::v1_to_v2_command(cmd, args);
proxy_daemon_send_v2(root, v2_cmd).await
}
pub(crate) async fn proxy_daemon_v2(
root: &Path,
cmd: super::protocol::Command,
) -> ProxyDaemonResult {
let v2_cmd = match serde_json::to_value(&cmd) {
Ok(v) => v,
Err(_) => return ProxyDaemonResult::Unresponsive,
};
let result = match proxy_daemon_send_v2(root, v2_cmd.clone()).await {
AttemptOutcome::Final(result) => result,
AttemptOutcome::Retryable(_) => {
tokio::time::sleep(Duration::from_millis(100)).await;
match proxy_daemon_send_v2(root, v2_cmd.clone()).await {
AttemptOutcome::Final(result) | AttemptOutcome::Retryable(result) => result,
}
}
};
if matches!(
&result,
ProxyDaemonResult::NotRunning | ProxyDaemonResult::StaleSocket
) && super::daemon_lifecycle::ensure_daemon(root).await
{
match proxy_daemon_send_v2(root, v2_cmd).await {
AttemptOutcome::Final(r) | AttemptOutcome::Retryable(r) => return r,
}
}
result
}
async fn proxy_daemon_send_v2(root: &Path, v2_cmd: serde_json::Value) -> AttemptOutcome {
let sock_path = root.join("mati.sock");
if sock_path.as_os_str().len() > UNIX_SOCK_PATH_MAX {
tracing::warn!(
path = %sock_path.display(),
"mcp proxy: socket path exceeds Unix limit"
);
return AttemptOutcome::Final(ProxyDaemonResult::NotRunning);
}
if !sock_path.exists() {
return AttemptOutcome::Retryable(ProxyDaemonResult::NotRunning);
}
let stream = match UnixStream::connect(&sock_path).await {
Ok(s) => s,
Err(e) => {
let is_refused = e.kind() == std::io::ErrorKind::ConnectionRefused;
if is_refused {
use super::metadata::{self as meta, StaleCheckResult};
match meta::check_and_cleanup_stale(root) {
StaleCheckResult::StaleRemoved | StaleCheckResult::Clean => {
return AttemptOutcome::Retryable(ProxyDaemonResult::StaleSocket);
}
StaleCheckResult::OrphanSocket => {
let _ = std::fs::remove_file(&sock_path);
return AttemptOutcome::Retryable(ProxyDaemonResult::StaleSocket);
}
StaleCheckResult::LiveDaemon { .. } => {
return AttemptOutcome::Retryable(ProxyDaemonResult::Unresponsive);
}
}
}
return AttemptOutcome::Retryable(ProxyDaemonResult::NotRunning);
}
};
let daemon_session = super::metadata::read_metadata(root)
.map(|m| m.session)
.unwrap_or_else(uuid::Uuid::nil);
let request = serde_json::json!({
"v": super::protocol::PROTOCOL_VERSION,
"id": uuid::Uuid::new_v4(),
"session": daemon_session,
"cmd": v2_cmd,
});
let (reader, mut writer) = stream.into_split();
let mut bytes = match serde_json::to_vec(&request) {
Ok(b) => b,
Err(_) => return AttemptOutcome::Final(ProxyDaemonResult::Unresponsive),
};
bytes.push(b'\n');
if writer.write_all(&bytes).await.is_err() {
return AttemptOutcome::Retryable(ProxyDaemonResult::Unresponsive);
}
if writer.shutdown().await.is_err() {
return AttemptOutcome::Retryable(ProxyDaemonResult::Unresponsive);
}
let mut buf_reader = BufReader::new(reader);
let mut line = String::new();
match tokio::time::timeout(Duration::from_secs(2), buf_reader.read_line(&mut line)).await {
Ok(Ok(n)) if n > 0 => {}
_ => return AttemptOutcome::Retryable(ProxyDaemonResult::Unresponsive),
}
let resp: serde_json::Value = match serde_json::from_str(line.trim()) {
Ok(v) => v,
Err(_) => return AttemptOutcome::Final(ProxyDaemonResult::Unresponsive),
};
match resp.get("status").and_then(|s| s.as_str()) {
Some("ok") => {
let data = resp.get("data").cloned().unwrap_or(serde_json::Value::Null);
AttemptOutcome::Final(ProxyDaemonResult::Ok(
serde_json::json!({"ok": true, "v": 2, "data": data}),
))
}
Some("err") => {
let code = resp
.get("code")
.and_then(|c| c.as_str())
.unwrap_or("internal");
let message = resp
.get("message")
.and_then(|m| m.as_str())
.unwrap_or("unknown error");
let envelope = serde_json::json!({
"ok": false, "v": 2, "error": message, "code": code
});
if code == "session_mismatch" {
tracing::debug!(
"mcp proxy: session mismatch — daemon may have restarted, will retry"
);
AttemptOutcome::Retryable(ProxyDaemonResult::Ok(envelope))
} else {
AttemptOutcome::Final(ProxyDaemonResult::Ok(envelope))
}
}
_ => AttemptOutcome::Retryable(ProxyDaemonResult::Unresponsive),
}
}