pitchfork-cli 2.26.0

Daemons with DX
Documentation
//! IPC request handling and dispatch
//!
//! Handles incoming IPC requests from CLI clients and routes them to the appropriate handlers.

use super::{SUPERVISOR, Supervisor};
use crate::Result;
use crate::ipc::server::IpcServer;
use crate::ipc::{IpcRequest, IpcResponse};
use miette::IntoDiagnostic;

const VERSION: &str = env!("CARGO_PKG_VERSION");

/// Dedupe the supervisor-side version-mismatch warning. Parallel batch clients
/// open one dedicated connection per daemon, each performing its own handshake,
/// which would otherwise log the warning once per connection. Keyed by client
/// version so a different mismatched client still gets logged.
fn version_mismatch_should_warn(client_version: &str) -> bool {
    static WARNED: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
    let mut warned = WARNED.lock().unwrap_or_else(|e| e.into_inner());
    if warned.as_deref() == Some(client_version) {
        false
    } else {
        *warned = Some(client_version.to_string());
        true
    }
}

impl Supervisor {
    /// Main IPC connection watch loop - reads and dispatches requests
    pub(crate) async fn conn_watch(&self, mut ipc: IpcServer) -> ! {
        loop {
            let (msg, send) = match ipc.read().await {
                Ok(msg) => msg,
                Err(e) => {
                    error!("failed to accept connection: {e:?}");
                    continue;
                }
            };
            debug!("received message: {msg:?}");
            tokio::spawn(async move {
                let rsp = SUPERVISOR
                    .handle_ipc(msg)
                    .await
                    .unwrap_or_else(|err| IpcResponse::Error(err.to_string()));
                if let Err(err) = send.send(rsp).await {
                    debug!("failed to send message: {err:?}");
                }
            });
        }
    }

    /// Handle a single IPC request and return the appropriate response
    pub(crate) async fn handle_ipc(&self, req: IpcRequest) -> Result<IpcResponse> {
        let rsp = match req {
            IpcRequest::Invalid { error } => {
                warn!("Invalid IPC request: {error}");
                return Ok(IpcResponse::Error(format!("Invalid request: {error}")));
            }
            IpcRequest::Connect => {
                debug!("received connect message (legacy, no version info)");
                IpcResponse::Ok
            }
            IpcRequest::ConnectV2 {
                version: client_version,
            } => {
                debug!("received connect message (client version: {client_version})");
                if client_version != VERSION && version_mismatch_should_warn(&client_version) {
                    warn!(
                        "Client version {client_version} differs from supervisor version {VERSION}. \
                            Restart the supervisor with: pitchfork supervisor start --force"
                    );
                }
                IpcResponse::ConnectOk {
                    version: VERSION.to_string(),
                }
            }
            IpcRequest::Stop { id } => {
                // id is already DaemonId, no validation needed
                self.stop(&id).await?
            }
            IpcRequest::Run(opts) => {
                // opts.id is already DaemonId, no validation needed
                self.run(opts).await?
            }
            IpcRequest::Enable { id } => {
                // id is already DaemonId, no validation needed
                if self.enable(&id).await? {
                    IpcResponse::Yes
                } else {
                    IpcResponse::No
                }
            }
            IpcRequest::Disable { id } => {
                // id is already DaemonId, no validation needed
                if self.disable(&id).await? {
                    IpcResponse::Yes
                } else {
                    IpcResponse::No
                }
            }
            IpcRequest::GetActiveDaemons => {
                let daemons = self.active_daemons().await;
                IpcResponse::ActiveDaemons(daemons)
            }
            IpcRequest::GetNotifications => {
                let notifications = self.get_notifications().await;
                IpcResponse::Notifications(notifications)
            }
            IpcRequest::UpdateShellDir { shell_pid, dir } => {
                let prev = self.get_shell_dir(shell_pid).await;
                self.set_shell_dir(shell_pid, dir.clone()).await?;
                // Cancel any pending autostops for daemons in the new directory
                self.cancel_pending_autostops_for_dir(&dir).await;
                if let Some(prev) = prev {
                    self.leave_dir(&prev).await?;
                }
                self.refresh().await?;
                IpcResponse::Ok
            }
            IpcRequest::SinkOutputLine {
                id,
                token,
                fires_hook,
                line,
            } => {
                super::log_sink::deliver_reported_line(&id, token, fires_hook, line).await;
                IpcResponse::Ok
            }
            IpcRequest::Clean => {
                self.clean().await?;
                IpcResponse::Ok
            }
            IpcRequest::CleanFiltered {
                namespaces,
                daemons,
                prune,
            } => {
                let count = self.clean_filtered(&namespaces, &daemons, prune).await?;
                IpcResponse::Cleaned { count }
            }
            IpcRequest::GetDisabledDaemons => {
                let disabled = self.state_file.lock().await.disabled.clone();
                IpcResponse::DisabledDaemons(disabled.into_iter().collect())
            }
            IpcRequest::SyncMdns => {
                self.sync_mdns().await;
                IpcResponse::MdnsSynced
            }
            IpcRequest::ReloadConfig => {
                tokio::task::spawn_blocking(|| {
                    crate::pitchfork_toml::invalidate_config_cache();
                    crate::settings::reload_settings();
                    crate::logger::apply_settings();
                })
                .await
                .into_diagnostic()?;
                IpcResponse::ConfigReloaded
            }
            IpcRequest::ProjectEnter { pid, dir } => {
                debug!("handling project enter pid {pid} dir {}", dir.display());
                let prev = self.enter_project_session(pid, dir.clone()).await?;
                self.cancel_pending_autostops_for_dir(&dir).await;
                // When re-entering (prev.is_some()), the new session keeps the
                // directory active, so leave_dir would be a no-op. Skip it to
                // avoid unnecessary autostop evaluation.
                let _ = prev;
                self.refresh().await?;
                IpcResponse::Ok
            }
            IpcRequest::ProjectLeave { pid, dir } => {
                debug!("handling project leave pid {pid} dir {}", dir.display());
                if let Some(left_dir) = self.leave_project_session(pid, &dir).await? {
                    debug!(
                        "project leave removed session pid {pid}, evaluating {left_dir:?} for autostop"
                    );
                    self.leave_dir(&left_dir).await?;
                } else {
                    debug!(
                        "project leave: session pid {pid} dir {} not found",
                        dir.display()
                    );
                }
                self.refresh().await?;
                IpcResponse::Ok
            }
            IpcRequest::GetProjectSessions => {
                let sessions = self.get_project_sessions_info().await;
                IpcResponse::ProjectSessions(sessions)
            }
            IpcRequest::GetWebUrl => IpcResponse::WebUrl {
                url: crate::web::url(),
            },
        };
        // Ensure state is flushed to disk before returning the response
        // so that CLI commands reading StateFile::get() see fresh data.
        self.flush_state().await;
        Ok(rsp)
    }
}