trusty-mpm 1.8.2

trusty-mpm: unified multi-agent orchestration platform (core, daemon, CLI, TUI, Telegram)
#![allow(dead_code)]
//! A stand-in trusty-* daemon on a temp Unix socket, for tests (#6286, #6285).
//!
//! Why: this crate's memory rigs each bound an ephemeral TCP port and served
//! `POST /rpc` (or the REST routes) with axum, because that is how they reached
//! trusty-memory. ADR-0032 retired that listener, so every one of them has to
//! dial a socket instead — and a copy of the accept loop per rig is the
//! duplication the workspace's common-entry-point rule exists to prevent.
//! #6285 moved the trusty-search rigs across the same way, which is why nothing
//! here names a service: the handler decides which daemon it is pretending to
//! be.
//!
//! What: [`spawn`] binds a socket under a `TempDir`, mounts `handler` as the
//! router's catch-all through the same [`trusty_common::uds::server`] pieces the
//! real daemon uses, and serves until the returned [`MockUdsDaemon`] drops.
//! [`spawn_at`] does the same at a path the caller chose, for a rig that has to
//! bind where production discovery will look.
//! The handler answers a `result` value directly — a rig that needs the
//! `tools/call` envelope wraps it itself, the same way it did over HTTP.
//!
//! This is a `#[cfg(test)]` module, so it never ships.
//!
//! Test: every caller — `core::memory_import::tests`,
//! `tui::coordinator::tests`, `daemon::doctor_tests`,
//! `daemon::doctor_search_pin_tests`,
//! `session_manager::search_gc_guard_tests`,
//! `session_manager::index_delete_guard::tests`.

use std::future::Future;
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::sync::Arc;

use async_trait::async_trait;
use serde_json::Value;
use tempfile::TempDir;
pub use trusty_common::uds::server::RpcError;
use trusty_common::uds::server::{RpcFallback, RpcRouter, RpcServeOptions, serve_until};

/// What one mock call answers, as a boxed future.
///
/// Boxed rather than generic because one rig's handler awaits a `watch`
/// channel: it cannot be a plain function of its arguments. `Err` is how a rig
/// makes the daemon refuse — the palace-missing case `ensure_palace` branches
/// on used a JSON-RPC error body over HTTP too.
pub type MockFuture = Pin<Box<dyn Future<Output = Result<Value, RpcError>> + Send>>;

/// A running mock daemon. Dropping it stops the accept loop and removes the
/// socket with its temp directory.
pub struct MockUdsDaemon {
    socket: PathBuf,
    /// Held only when [`spawn`] minted the directory. [`spawn_at`] binds inside
    /// a directory the caller owns and keeps alive.
    _dir: Option<TempDir>,
    shutdown: Option<tokio::sync::oneshot::Sender<()>>,
}

impl MockUdsDaemon {
    /// The path a client under test should dial.
    pub fn socket(&self) -> &Path {
        &self.socket
    }
}

impl Drop for MockUdsDaemon {
    fn drop(&mut self) {
        if let Some(tx) = self.shutdown.take() {
            let _ = tx.send(());
        }
    }
}

/// The catch-all that hands every method to the test's closure.
struct MockFallback<F> {
    handler: F,
}

#[async_trait]
impl<F> RpcFallback for MockFallback<F>
where
    F: Fn(&str, Value) -> MockFuture + Send + Sync + 'static,
{
    async fn call(&self, method: &str, params: Value) -> Result<Value, RpcError> {
        (self.handler)(method, params).await
    }
}

/// Start a mock daemon answering every method through `handler`.
///
/// # Panics
///
/// When the socket cannot be bound — a test-only failure with no recovery.
pub async fn spawn<F>(handler: F) -> MockUdsDaemon
where
    F: Fn(&str, Value) -> MockFuture + Send + Sync + 'static,
{
    let dir = TempDir::new().expect("tempdir for the mock socket");
    let socket = dir.path().join("daemon.sock");
    let mut daemon = spawn_at(socket, handler).await;
    daemon._dir = Some(dir);
    daemon
}

/// Start a mock daemon at a socket path the CALLER chose.
///
/// Why: a rig that exercises production socket DISCOVERY cannot take whatever
/// path [`spawn`] minted — it has to bind where `daemon_socket_path` will look,
/// under a `TRUSTY_DATA_DIR_OVERRIDE` the test controls. `search_gc_guard_tests`
/// and `index_delete_guard_tests` both need that, and a second accept loop for
/// them is the duplication this module exists to prevent (#6285).
/// What: as [`spawn`], except the temp directory keeping the socket alive is the
/// caller's. `bind_hardened` creates and hardens the parent directory itself.
///
/// # Panics
///
/// When the socket cannot be bound — a test-only failure with no recovery.
pub async fn spawn_at<F>(socket: PathBuf, handler: F) -> MockUdsDaemon
where
    F: Fn(&str, Value) -> MockFuture + Send + Sync + 'static,
{
    let listener = trusty_common::uds::bind_hardened(&socket).expect("bind the mock socket");

    let router = Arc::new(RpcRouter::new().fallback(MockFallback { handler }));
    let (tx, rx) = tokio::sync::oneshot::channel::<()>();
    tokio::spawn(async move {
        serve_until(&listener, router, RpcServeOptions::default(), async {
            let _ = rx.await;
        })
        .await;
    });

    MockUdsDaemon {
        socket,
        _dir: None,
        shutdown: Some(tx),
    }
}

/// Wrap `inner` the way the daemon's `tools/call` arm answers.
///
/// Why: the real dispatcher stringifies a tool's result into
/// `result.content[0].text`, and the rigs assert the unwrap as well as the
/// request. Keeping the shape here means one place gets it wrong or right.
pub fn tools_call_envelope(inner: &Value) -> Value {
    serde_json::json!({ "content": [{ "type": "text", "text": inner.to_string() }] })
}

/// Sugar for a handler that answers the same value to every call.
pub fn always(result: Value) -> impl Fn(&str, Value) -> MockFuture + Send + Sync + 'static {
    move |_method, _params| {
        let result = result.clone();
        Box::pin(async move { Ok(result) })
    }
}

/// A mock daemon owning its own runtime on its own OS thread (#7237).
///
/// Why: the `search_index` registration rigs are SYNCHRONOUS `#[test]`s holding
/// process-global env guards, and the code they drive is a blocking function.
/// An async rig would have to hold those guards across an `.await`. Giving the
/// daemon its own thread keeps the rig exactly the shape the retired
/// `TcpListener` fixtures had.
/// What: dropping it stops the accept loop and joins the thread.
pub struct BlockingMockDaemon {
    socket: PathBuf,
    shutdown: Option<tokio::sync::oneshot::Sender<()>>,
    thread: Option<std::thread::JoinHandle<()>>,
}

impl BlockingMockDaemon {
    /// The path a client under test should dial.
    pub fn socket(&self) -> &Path {
        &self.socket
    }
}

impl Drop for BlockingMockDaemon {
    fn drop(&mut self) {
        if let Some(tx) = self.shutdown.take() {
            let _ = tx.send(());
        }
        if let Some(thread) = self.thread.take() {
            let _ = thread.join();
        }
    }
}

/// Start a [`BlockingMockDaemon`] at `socket`, returning only once it is bound.
///
/// Why the readiness handshake: a rig that returned before `bind_hardened` ran
/// would race the client's dial against the bind and flake.
///
/// # Panics
///
/// When the runtime cannot be built or the socket cannot be bound — test-only
/// failures with no recovery.
pub fn spawn_blocking_at<F>(socket: PathBuf, handler: F) -> BlockingMockDaemon
where
    F: Fn(&str, Value) -> MockFuture + Send + Sync + 'static,
{
    let (ready_tx, ready_rx) = std::sync::mpsc::channel::<()>();
    let (stop_tx, stop_rx) = tokio::sync::oneshot::channel::<()>();
    let bind_at = socket.clone();

    let thread = std::thread::spawn(move || {
        let runtime = tokio::runtime::Builder::new_current_thread()
            .enable_all()
            .build()
            .expect("runtime for the mock daemon");
        runtime.block_on(async move {
            let listener =
                trusty_common::uds::bind_hardened(&bind_at).expect("bind the mock socket");
            let router = Arc::new(RpcRouter::new().fallback(MockFallback { handler }));
            let _ = ready_tx.send(());
            serve_until(&listener, router, RpcServeOptions::default(), async {
                let _ = stop_rx.await;
            })
            .await;
        });
    });

    ready_rx
        .recv()
        .expect("the mock daemon must bind before the rig proceeds");
    BlockingMockDaemon {
        socket,
        shutdown: Some(stop_tx),
        thread: Some(thread),
    }
}