scryer-mcp 0.2.1

Model Context Protocol (MCP) server for Scryer code intelligence
//! Thin-client side of the daemon: connect (spawning the daemon if needed) and pipe stdio.

use std::fs::File;
use std::path::Path;
use std::process::{Command, Stdio};
use std::time::{Duration, Instant};

use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, BufReader};
use tokio::net::UnixStream;
use tokio::net::unix::{OwnedReadHalf, OwnedWriteHalf};

use super::handshake::{
    ControlCommand, Hello, HelloMode, HelloReply, read_json_line, write_json_line,
};
use super::paths::DaemonPaths;
use super::prepare_run_dir;
use super::server::daemon_running;

/// How long to wait for the daemon's handshake reply.
const REPLY_TIMEOUT: Duration = Duration::from_secs(30);

/// Log size above which the spawner truncates the daemon log.
const MAX_LOG_BYTES: u64 = 5 * 1024 * 1024;

/// Lines of the daemon log shown when the daemon fails to start.
const LOG_TAIL_LINES: usize = 20;

/// A connection that has completed the handshake; raw MCP traffic flows from here.
pub struct DaemonConnection {
    pub reader: BufReader<OwnedReadHalf>,
    pub writer: OwnedWriteHalf,
    pub reply: HelloReply,
}

/// Connect to a running daemon, if there is one.
pub async fn connect(paths: &DaemonPaths) -> std::io::Result<UnixStream> {
    if paths.socket_is_external() {
        verify_socket_owner(paths)?;
    }
    UnixStream::connect(&paths.socket).await
}

/// Refuse a socket in a shared directory (e.g. `/tmp`) that another user created.
fn verify_socket_owner(paths: &DaemonPaths) -> std::io::Result<()> {
    use std::os::unix::fs::MetadataExt;

    let socket_uid = std::fs::symlink_metadata(&paths.socket)?.uid();
    let own_uid = std::fs::metadata(paths.run_dir())?.uid();
    if socket_uid != own_uid {
        return Err(std::io::Error::new(
            std::io::ErrorKind::PermissionDenied,
            format!(
                "{} is owned by another user; refusing to connect",
                paths.socket.display()
            ),
        ));
    }
    Ok(())
}

/// Connect to the daemon for `paths`, calling `spawn` to start it if none is running.
///
/// Startup is serialized by the spawn lock, so concurrent clients start exactly one daemon.
pub async fn connect_or_spawn<F>(
    paths: &DaemonPaths,
    spawn: F,
    startup_timeout: Duration,
) -> anyhow::Result<UnixStream>
where
    F: FnOnce(&DaemonPaths) -> anyhow::Result<()>,
{
    if let Ok(stream) = connect(paths).await {
        return Ok(stream);
    }

    prepare_run_dir(paths.run_dir())?;
    let lock_path = paths.spawn_lock.clone();
    let _spawn_lock = tokio::task::spawn_blocking(move || -> std::io::Result<File> {
        let file = File::options()
            .write(true)
            .create(true)
            .truncate(false)
            .open(lock_path)?;
        file.lock()?;
        Ok(file)
    })
    .await??;

    // Another client may have started the daemon while we waited for the lock.
    if let Ok(stream) = connect(paths).await {
        return Ok(stream);
    }

    // A live daemon without a socket is still starting up; only spawn when none holds the lock.
    if !daemon_running(paths) {
        let _ = std::fs::remove_file(&paths.socket);
        spawn(paths)?;
    }

    let deadline = Instant::now() + startup_timeout;
    loop {
        match connect(paths).await {
            Ok(stream) => return Ok(stream),
            Err(_) if Instant::now() < deadline => {
                tokio::time::sleep(Duration::from_millis(50)).await;
            }
            Err(e) => anyhow::bail!(
                "Scryer daemon did not accept connections on {} within {startup_timeout:?} ({e}).{}",
                paths.socket.display(),
                log_tail(&paths.log_file)
            ),
        }
    }
}

/// How long [`replace_stale_daemon`] waits for a stale daemon to exit.
const STALE_STOP_TIMEOUT: Duration = Duration::from_secs(10);

/// Whether a daemon replying with `build` runs a different binary than this process.
///
/// A daemon that sends no build ID predates build checks, so it is a different binary too.
pub fn is_stale_build(build: Option<&str>) -> bool {
    build != Some(crate::build_info::build_id())
}

/// Stop the running daemon if it was started from a different binary and has no sessions.
///
/// The daemon outlives its clients, so after a rebuild or reinstall the next client would
/// otherwise keep talking to old code. A daemon still serving sessions is left running.
/// Returns `true` if a stale daemon was stopped.
pub async fn replace_stale_daemon(paths: &DaemonPaths) -> anyhow::Result<bool> {
    let Some(reply) = control(paths, ControlCommand::Status).await? else {
        return Ok(false);
    };
    let Some(status) = reply.status else {
        return Ok(false);
    };
    if !is_stale_build(status.build.as_deref()) || status.sessions > 0 {
        return Ok(false);
    }

    tracing::info!(
        "Stopping Scryer daemon pid {} (build {}): this client is build {}",
        status.pid,
        status.build.as_deref().unwrap_or("unknown"),
        crate::build_info::build_id()
    );
    // Stop it only if it is still idle: a session may have connected since `Status`. A daemon
    // too old to know `stop_if_idle` rejects it, so fall back to a plain stop for those.
    let stopped = match control(paths, ControlCommand::StopIfIdle).await {
        Ok(Some(reply)) => reply.status.is_none(),
        Ok(None) => false,
        Err(_) => matches!(control(paths, ControlCommand::Stop).await, Ok(Some(_))),
    };
    if !stopped {
        tracing::info!(
            "Stale Scryer daemon pid {} gained a session; leaving it",
            status.pid
        );
        return Ok(false);
    }
    // Wait for the stale daemon itself to exit. Another client may start a fresh daemon the
    // moment the lock frees, so "no daemon running" alone could never become true.
    let deadline = Instant::now() + STALE_STOP_TIMEOUT;
    while daemon_running(paths) && pid_file_owner(paths) == Some(status.pid) {
        if Instant::now() >= deadline {
            anyhow::bail!(
                "stale Scryer daemon (pid {}) did not exit within {STALE_STOP_TIMEOUT:?}",
                status.pid
            );
        }
        tokio::time::sleep(Duration::from_millis(50)).await;
    }
    Ok(true)
}

/// PID recorded in the daemon's pid file, if it can be read.
fn pid_file_owner(paths: &DaemonPaths) -> Option<u32> {
    std::fs::read_to_string(&paths.pid_file)
        .ok()?
        .trim()
        .parse()
        .ok()
}

/// Send `hello` and wait for the daemon's reply.
pub async fn handshake(stream: UnixStream, hello: &Hello) -> anyhow::Result<DaemonConnection> {
    let (read, mut writer) = stream.into_split();
    write_json_line(&mut writer, hello).await?;

    let mut reader = BufReader::new(read);
    let reply: HelloReply = tokio::time::timeout(REPLY_TIMEOUT, read_json_line(&mut reader))
        .await
        .map_err(|_| anyhow::anyhow!("timed out waiting for the Scryer daemon handshake"))??;
    if !reply.ok {
        anyhow::bail!(
            "Scryer daemon rejected the connection: {}",
            reply.error.as_deref().unwrap_or("unknown error")
        );
    }

    Ok(DaemonConnection {
        reader,
        writer,
        reply,
    })
}

/// Send a control command to a running daemon. Returns `None` if no daemon is listening.
pub async fn control(
    paths: &DaemonPaths,
    cmd: ControlCommand,
) -> anyhow::Result<Option<HelloReply>> {
    let Ok(stream) = connect(paths).await else {
        return Ok(None);
    };
    let conn = handshake(stream, &Hello::new(HelloMode::Control { cmd })).await?;
    Ok(Some(conn.reply))
}

/// Start `exe daemon run` for `paths` as a detached background process.
///
/// The daemon's stderr goes to `paths.log_file`; stdin and stdout are closed.
pub fn spawn_daemon_process(
    exe: &Path,
    paths: &DaemonPaths,
    idle_timeout_secs: u64,
) -> anyhow::Result<()> {
    use std::os::unix::process::CommandExt;

    prepare_run_dir(paths.run_dir())?;
    if std::fs::metadata(&paths.log_file).is_ok_and(|m| m.len() > MAX_LOG_BYTES) {
        let _ = std::fs::remove_file(&paths.log_file);
    }
    let log = File::options()
        .create(true)
        .append(true)
        .open(&paths.log_file)?;

    let mut child = Command::new(exe)
        .arg("daemon")
        .arg("run")
        .arg("--db-url")
        .arg(&paths.db_path)
        .arg("--idle-timeout")
        .arg(idle_timeout_secs.to_string())
        .env_remove("SCRYER_DB_URL")
        .stdin(Stdio::null())
        .stdout(Stdio::null())
        .stderr(log)
        // New process group: the editor killing the client must not take the daemon with it.
        .process_group(0)
        .spawn()?;

    // Reap the daemon if it exits while this client is still alive.
    std::thread::spawn(move || {
        let _ = child.wait();
    });
    Ok(())
}

/// Copy `input` to the daemon and the daemon's output to `output` until the session ends.
///
/// When `input` closes, the socket's write half is shut down and remaining daemon output is
/// drained. When the daemon closes the socket, the pipe ends immediately.
pub async fn pipe<I, O, R, W>(
    mut input: I,
    mut output: O,
    mut daemon_read: R,
    mut daemon_write: W,
) -> std::io::Result<()>
where
    I: AsyncRead + Unpin,
    O: AsyncWrite + Unpin,
    R: AsyncRead + Unpin,
    W: AsyncWrite + Unpin,
{
    let upstream = async {
        let copied = tokio::io::copy(&mut input, &mut daemon_write).await;
        let _ = daemon_write.shutdown().await;
        copied.map(|_| ())
    };
    let downstream = async {
        let mut buf = vec![0u8; 64 * 1024];
        loop {
            let n = daemon_read.read(&mut buf).await?;
            if n == 0 {
                return Ok::<_, std::io::Error>(());
            }
            output.write_all(&buf[..n]).await?;
            output.flush().await?;
        }
    };
    tokio::pin!(upstream, downstream);

    tokio::select! {
        res = &mut downstream => res,
        res = &mut upstream => {
            res?;
            downstream.await
        }
    }
}

/// Format the last lines of the daemon log for an error message.
fn log_tail(log_file: &Path) -> String {
    let Ok(content) = std::fs::read_to_string(log_file) else {
        return String::new();
    };
    let lines: Vec<&str> = content.lines().collect();
    if lines.is_empty() {
        return String::new();
    }
    let tail = &lines[lines.len().saturating_sub(LOG_TAIL_LINES)..];
    format!(
        "\nLast lines of {}:\n{}",
        log_file.display(),
        tail.join("\n")
    )
}