Skip to main content

pitchfork_cli/ipc/
mod.rs

1use crate::Result;
2use crate::daemon::{Daemon, RunOptions};
3use crate::daemon_id::DaemonId;
4use crate::env;
5use interprocess::local_socket::Name;
6#[cfg(unix)]
7use interprocess::local_socket::{GenericFilePath, ToFsName};
8#[cfg(windows)]
9use interprocess::local_socket::{GenericNamespaced, ToNsName};
10use miette::{Context, IntoDiagnostic};
11use std::path::PathBuf;
12
13pub(crate) mod batch;
14pub(crate) mod client;
15pub(crate) mod server;
16
17// #[derive(Debug, Clone, serde::Serialize, serde::Deserialize, strum::Display, strum::EnumIs)]
18// pub enum IpcMessage {
19//     Connect(String),
20//     ConnectOK,
21//     Run(String, Vec<String>),
22//     Stop(String),
23//     DaemonAlreadyRunning(String),
24//     DaemonAlreadyStopped(String),
25//     DaemonStart(Daemon),
26//     DaemonStop { name: String },
27//     DaemonFailed { name: String, error: String },
28//     Response(String),
29// }
30
31#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, strum::Display, strum::EnumIs)]
32#[allow(clippy::large_enum_variant)]
33pub enum IpcRequest {
34    Connect,
35    /// Versioned connect handshake (v2): client sends its version so the supervisor can
36    /// detect mismatches. Kept as a separate variant so the wire format of `Connect`
37    /// (unit variant) stays unchanged for backward compatibility with older supervisors.
38    ConnectV2 {
39        version: String,
40    },
41    Clean,
42    Stop {
43        id: DaemonId,
44    },
45    GetActiveDaemons,
46    GetDisabledDaemons,
47    Run(RunOptions),
48    Enable {
49        id: DaemonId,
50    },
51    Disable {
52        id: DaemonId,
53    },
54    UpdateShellDir {
55        shell_pid: u32,
56        dir: PathBuf,
57    },
58    GetNotifications,
59    /// Notify the supervisor that the slug registry has changed (e.g. `proxy add/remove`).
60    /// The supervisor should re-read slugs and update mDNS records accordingly.
61    SyncMdns,
62    /// Notify the supervisor that settings have changed.
63    /// The supervisor should reload settings from config files.
64    ReloadConfig,
65    /// Enter or replace a project session for a host PID in a directory.
66    ProjectEnter {
67        pid: u32,
68        dir: PathBuf,
69    },
70    /// Leave a project session for a host PID in a directory.
71    ProjectLeave {
72        pid: u32,
73        dir: PathBuf,
74    },
75    /// List all tracked project sessions with live liveness status filled in
76    /// by the supervisor.
77    GetProjectSessions,
78    /// Invalid request (failed to deserialize)
79    #[serde(skip)]
80    Invalid {
81        error: String,
82    },
83}
84
85/// A snapshot of a single project session, returned by `GetProjectSessions`.
86///
87/// `liveness_title` is the title recorded at enter time. `alive` and
88/// `current_title` are filled in by the supervisor from its `PROCS` singleton
89/// so the client does not need process introspection.
90#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
91pub struct ProjectSessionInfo {
92    pub pid: u32,
93    pub directory: PathBuf,
94    #[serde(skip_serializing_if = "Option::is_none", default)]
95    pub liveness_title: Option<String>,
96    pub alive: bool,
97    #[serde(skip_serializing_if = "Option::is_none", default)]
98    pub current_title: Option<String>,
99}
100
101#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, strum::Display, strum::EnumIs)]
102pub enum IpcResponse {
103    Ok,
104    /// Successful connect handshake, includes supervisor version for mismatch detection
105    ConnectOk {
106        version: String,
107    },
108    Yes,
109    No,
110    Error(String),
111    Notifications(Vec<(log::LevelFilter, String)>),
112    ActiveDaemons(Vec<Daemon>),
113    DisabledDaemons(Vec<DaemonId>),
114    DaemonAlreadyRunning,
115    DaemonStart {
116        daemon: Daemon,
117    },
118    DaemonFailed {
119        error: String,
120    },
121    /// Port conflict detected with detailed process information
122    PortConflict {
123        port: u16,
124        process: String,
125        pid: u32,
126    },
127    /// No available ports found after exhausting auto-bump attempts
128    NoAvailablePort {
129        start_port: u16,
130        attempts: u32,
131    },
132    DaemonReady {
133        daemon: Daemon,
134    },
135    DaemonFailedWithCode {
136        exit_code: Option<i32>,
137    },
138    /// Process was not running but had a PID record (unexpected exit)
139    DaemonWasNotRunning,
140    /// mDNS sync completed (or was a no-op if LAN mode is disabled)
141    MdnsSynced,
142    /// Settings reloaded from config files
143    ConfigReloaded,
144    /// Failed to kill the process (still running)
145    DaemonStopFailed {
146        error: String,
147    },
148    /// Daemon exists but is not running (no PID)
149    DaemonNotRunning,
150    DaemonNotFound,
151    /// Snapshot of all project sessions (response to `GetProjectSessions`).
152    ProjectSessions(Vec<ProjectSessionInfo>),
153}
154fn fs_name(name: &str) -> Result<Name<'_>> {
155    // Unix: use a filesystem path for the AF_UNIX socket.
156    #[cfg(unix)]
157    {
158        let path = env::IPC_SOCK_DIR.join(name).with_extension("sock");
159        let fs_name = path.to_fs_name::<GenericFilePath>().into_diagnostic()?;
160        Ok(fs_name)
161    }
162    // Windows: named pipes use a flat namespace (\\.\pipe\<name>) that
163    // cannot contain path separators. Derive a unique pipe name from the
164    // state directory to preserve test isolation when multiple supervisors
165    // run concurrently with different PITCHFORK_STATE_DIR values.
166    //
167    // Use a hash of the state directory path rather than character replacement
168    // to guarantee injectivity: `C:\a.b` and `C:\a\b` would both flatten to
169    // `C--a-b` with the old approach, causing pipe name collisions.
170    #[cfg(windows)]
171    {
172        let state_dir = env::PITCHFORK_STATE_DIR.to_string_lossy();
173        // FNV-1a hash: deterministic, stable across Rust versions.
174        // DefaultHasher's algorithm is not guaranteed stable, which would
175        // break IPC if the CLI and supervisor were ever compiled with
176        // different toolchains.
177        let mut hash: u64 = 0xcbf29ce484222325;
178        for byte in state_dir.bytes() {
179            hash ^= byte as u64;
180            hash = hash.wrapping_mul(0x100000001b3);
181        }
182        let pipe_name = format!("pitchfork-{hash:016x}-{name}");
183        Ok(pipe_name
184            .to_ns_name::<GenericNamespaced>()
185            .into_diagnostic()?)
186    }
187}
188
189fn serialize<T: serde::Serialize>(msg: &T) -> Result<Vec<u8>> {
190    if *env::IPC_JSON {
191        serde_json::to_vec(msg)
192            .into_diagnostic()
193            .wrap_err("failed to serialize IPC message as JSON")
194    } else {
195        rmp_serde::to_vec(msg)
196            .into_diagnostic()
197            .wrap_err("failed to serialize IPC message as MessagePack")
198    }
199}
200
201fn deserialize<T: serde::de::DeserializeOwned>(bytes: &[u8]) -> Result<T> {
202    let mut bytes = bytes.to_vec();
203    bytes.pop();
204    let preview = std::str::from_utf8(&bytes).unwrap_or("<binary>");
205    trace!("msg: {preview:?}");
206    if *env::IPC_JSON {
207        serde_json::from_slice(&bytes)
208            .into_diagnostic()
209            .wrap_err("failed to deserialize IPC JSON response")
210    } else {
211        rmp_serde::from_slice(&bytes)
212            .into_diagnostic()
213            .wrap_err("failed to deserialize IPC MessagePack response")
214    }
215}