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    /// A daemon's log sink reporting a line the supervisor needs to act on:
79    /// one matching the readiness pattern, one that should fire the
80    /// `on_output` hook, or both.
81    ///
82    /// The sink owns the daemon's output stream, so it does the matching and
83    /// tells the supervisor rather than the other way around. Sent by
84    /// `pitchfork log-sink`, not by any user-facing command.
85    SinkOutputLine {
86        id: DaemonId,
87        /// Identifies the start attempt this sink belongs to, so a report from
88        /// a sink still draining a previous attempt cannot satisfy a retry.
89        token: u64,
90        /// Whether this line passed the `on_output` hook's filter and debounce.
91        /// A line reported only because it matched the readiness pattern must
92        /// not fire a hook that filters for something else.
93        fires_hook: bool,
94        line: String,
95    },
96    /// Ask the supervisor for the URL of the web UI, if it is running.
97    /// Reflects the actual bound address, not static config.
98    GetWebUrl,
99    /// Remove stopped daemon registrations matching the supplied filters.
100    /// Appended to preserve the wire indexes of existing variants.
101    CleanFiltered {
102        namespaces: Vec<String>,
103        daemons: Vec<DaemonId>,
104        prune: bool,
105    },
106    /// These daemons are being started explicitly: none of them is to be
107    /// stopped for inactivity from now on, even if the proxy started it.
108    /// Sent before the start itself, and answered once any idle stop already
109    /// under way for one of them has finished. Appended to preserve the wire
110    /// indexes of existing variants.
111    ClaimDaemons {
112        ids: Vec<DaemonId>,
113    },
114    /// Invalid request (failed to deserialize)
115    #[serde(skip)]
116    Invalid {
117        error: String,
118    },
119}
120
121/// A snapshot of a single project session, returned by `GetProjectSessions`.
122///
123/// `liveness_title` is the title recorded at enter time. `alive` and
124/// `current_title` are filled in by the supervisor from its `PROCS` singleton
125/// so the client does not need process introspection.
126#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
127pub struct ProjectSessionInfo {
128    pub pid: u32,
129    pub directory: PathBuf,
130    #[serde(skip_serializing_if = "Option::is_none", default)]
131    pub liveness_title: Option<String>,
132    pub alive: bool,
133    #[serde(skip_serializing_if = "Option::is_none", default)]
134    pub current_title: Option<String>,
135}
136
137#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, strum::Display, strum::EnumIs)]
138pub enum IpcResponse {
139    Ok,
140    /// Successful connect handshake, includes supervisor version for mismatch detection
141    ConnectOk {
142        version: String,
143    },
144    Yes,
145    No,
146    Error(String),
147    Notifications(Vec<(log::LevelFilter, String)>),
148    ActiveDaemons(Vec<Daemon>),
149    DisabledDaemons(Vec<DaemonId>),
150    DaemonAlreadyRunning,
151    DaemonStart {
152        daemon: Daemon,
153    },
154    DaemonFailed {
155        error: String,
156    },
157    /// Port conflict detected with detailed process information
158    PortConflict {
159        port: u16,
160        process: String,
161        pid: u32,
162    },
163    /// No available ports found after exhausting auto-bump attempts
164    NoAvailablePort {
165        start_port: u16,
166        attempts: u32,
167    },
168    DaemonReady {
169        daemon: Daemon,
170    },
171    DaemonFailedWithCode {
172        exit_code: Option<i32>,
173        /// Ports resolved by the failed attempt, so in-process retry hooks
174        /// observe the attempt's actual (post-bump) ports.
175        #[serde(default)]
176        resolved_ports: Vec<u16>,
177    },
178    /// Process was not running but had a PID record (unexpected exit)
179    DaemonWasNotRunning,
180    /// mDNS sync completed (or was a no-op if LAN mode is disabled)
181    MdnsSynced,
182    /// Settings reloaded from config files
183    ConfigReloaded,
184    /// URL of the running web UI, or `None` if the web UI is not running
185    WebUrl {
186        url: Option<String>,
187    },
188    /// Failed to kill the process (still running)
189    DaemonStopFailed {
190        error: String,
191    },
192    /// Daemon exists but is not running (no PID)
193    DaemonNotRunning,
194    DaemonNotFound,
195    /// Snapshot of all project sessions (response to `GetProjectSessions`).
196    ProjectSessions(Vec<ProjectSessionInfo>),
197    /// Number of daemon registrations removed by `CleanFiltered`.
198    Cleaned {
199        count: u64,
200    },
201}
202fn fs_name(name: &str) -> Result<Name<'_>> {
203    // Unix: use a filesystem path for the AF_UNIX socket.
204    #[cfg(unix)]
205    {
206        let path = env::IPC_SOCK_DIR.join(name).with_extension("sock");
207        let fs_name = path.to_fs_name::<GenericFilePath>().into_diagnostic()?;
208        Ok(fs_name)
209    }
210    // Windows: named pipes use a flat namespace (\\.\pipe\<name>) that
211    // cannot contain path separators. Derive a unique pipe name from the
212    // state directory to preserve test isolation when multiple supervisors
213    // run concurrently with different PITCHFORK_STATE_DIR values.
214    //
215    // Use a hash of the state directory path rather than character replacement
216    // to guarantee injectivity: `C:\a.b` and `C:\a\b` would both flatten to
217    // `C--a-b` with the old approach, causing pipe name collisions.
218    #[cfg(windows)]
219    {
220        let state_dir = env::PITCHFORK_STATE_DIR.to_string_lossy();
221        // FNV-1a hash: deterministic, stable across Rust versions.
222        // DefaultHasher's algorithm is not guaranteed stable, which would
223        // break IPC if the CLI and supervisor were ever compiled with
224        // different toolchains.
225        let mut hash: u64 = 0xcbf29ce484222325;
226        for byte in state_dir.bytes() {
227            hash ^= byte as u64;
228            hash = hash.wrapping_mul(0x100000001b3);
229        }
230        let pipe_name = format!("pitchfork-{hash:016x}-{name}");
231        pipe_name
232            .to_ns_name::<GenericNamespaced>()
233            .into_diagnostic()
234    }
235}
236
237/// Human-readable location of the supervisor's IPC endpoint, for messages.
238pub(crate) fn socket_display() -> String {
239    #[cfg(unix)]
240    {
241        env::IPC_SOCK_MAIN.display().to_string()
242    }
243    #[cfg(windows)]
244    {
245        "the supervisor named pipe".to_string()
246    }
247}
248
249/// Whether a supervisor is accepting connections on the IPC socket right now.
250///
251/// This is the ground truth for "is a supervisor running": the state-file
252/// record can be lost or rewritten while the supervisor keeps serving, and a
253/// supervisor started on that basis would take the socket over and leave the
254/// original running but unreachable. A stale socket file left by a crashed
255/// supervisor refuses connections, so it does not count.
256pub(crate) async fn supervisor_listening() -> bool {
257    use interprocess::local_socket::traits::tokio::Stream as _;
258    let Ok(name) = fs_name("main") else {
259        return false;
260    };
261    let connect = interprocess::local_socket::tokio::Stream::connect(name);
262    // A live supervisor accepts at once (the kernel queues the connection);
263    // the timeout only guards against a platform where connecting blocks.
264    match tokio::time::timeout(std::time::Duration::from_secs(1), connect).await {
265        Ok(Ok(_)) => true,
266        Ok(Err(err)) => {
267            trace!("no supervisor listening on the IPC socket: {err}");
268            false
269        }
270        Err(_) => {
271            debug!("timed out probing the IPC socket; treating it as not listening");
272            false
273        }
274    }
275}
276
277fn serialize<T: serde::Serialize>(msg: &T) -> Result<Vec<u8>> {
278    if *env::IPC_JSON {
279        serde_json::to_vec(msg)
280            .into_diagnostic()
281            .wrap_err("failed to serialize IPC message as JSON")
282    } else {
283        rmp_serde::to_vec(msg)
284            .into_diagnostic()
285            .wrap_err("failed to serialize IPC message as MessagePack")
286    }
287}
288
289fn deserialize<T: serde::de::DeserializeOwned>(bytes: &[u8]) -> Result<T> {
290    let mut bytes = bytes.to_vec();
291    bytes.pop();
292    let preview = std::str::from_utf8(&bytes).unwrap_or("<binary>");
293    trace!("msg: {preview:?}");
294    if *env::IPC_JSON {
295        serde_json::from_slice(&bytes)
296            .into_diagnostic()
297            .wrap_err("failed to deserialize IPC JSON response")
298    } else {
299        rmp_serde::from_slice(&bytes)
300            .into_diagnostic()
301            .wrap_err("failed to deserialize IPC MessagePack response")
302    }
303}
304
305#[cfg(test)]
306mod tests {
307    use super::*;
308
309    #[test]
310    fn filtered_clean_ipc_round_trips() {
311        let request = IpcRequest::CleanFiltered {
312            namespaces: vec!["worktree".to_string()],
313            daemons: vec![DaemonId::new("worktree", "api")],
314            prune: true,
315        };
316        let mut bytes = serialize(&request).unwrap();
317        bytes.push(b'\n');
318        let decoded: IpcRequest = deserialize(&bytes).unwrap();
319        match decoded {
320            IpcRequest::CleanFiltered {
321                namespaces,
322                daemons,
323                prune,
324            } => {
325                assert_eq!(namespaces, ["worktree"]);
326                assert_eq!(daemons, [DaemonId::new("worktree", "api")]);
327                assert!(prune);
328            }
329            other => panic!("unexpected request: {other:?}"),
330        }
331
332        let mut bytes = serialize(&IpcResponse::Cleaned { count: 3 }).unwrap();
333        bytes.push(b'\n');
334        let decoded: IpcResponse = deserialize(&bytes).unwrap();
335        assert!(matches!(decoded, IpcResponse::Cleaned { count: 3 }));
336    }
337
338    fn round_trip<T: serde::Serialize + serde::de::DeserializeOwned>(value: &T) -> T {
339        let mut bytes = serialize(value).unwrap();
340        bytes.push(b'\n');
341        deserialize(&bytes).unwrap()
342    }
343
344    #[test]
345    fn claim_daemons_ipc_round_trips() {
346        let ids = vec![DaemonId::new("proj", "api"), DaemonId::new("proj", "db")];
347        match round_trip(&IpcRequest::ClaimDaemons { ids: ids.clone() }) {
348            IpcRequest::ClaimDaemons { ids: decoded } => assert_eq!(decoded, ids),
349            other => panic!("unexpected request: {other:?}"),
350        }
351    }
352
353    /// The idle-shutdown ownership rides at the end of both positionally
354    /// encoded structs, after fields that are skipped when empty.
355    #[test]
356    fn proxy_idle_timeout_survives_the_ipc_encoding() {
357        let daemon = Daemon {
358            id: DaemonId::new("proj", "api"),
359            proxy_idle_timeout_ms: Some(900_000),
360            ..Default::default()
361        };
362        match round_trip(&IpcResponse::ActiveDaemons(vec![daemon])) {
363            IpcResponse::ActiveDaemons(daemons) => {
364                assert_eq!(daemons[0].proxy_idle_timeout_ms, Some(900_000));
365                assert!(!daemons[0].oneshot);
366            }
367            other => panic!("unexpected response: {other:?}"),
368        }
369
370        let opts = RunOptions {
371            id: DaemonId::new("proj", "api"),
372            proxy_idle_timeout_ms: Some(900_000),
373            ..Default::default()
374        };
375        match round_trip(&IpcRequest::Run(opts)) {
376            IpcRequest::Run(opts) => assert_eq!(opts.proxy_idle_timeout_ms, Some(900_000)),
377            other => panic!("unexpected request: {other:?}"),
378        }
379    }
380}