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 /// Invalid request (failed to deserialize)
107 #[serde(skip)]
108 Invalid {
109 error: String,
110 },
111}
112
113/// A snapshot of a single project session, returned by `GetProjectSessions`.
114///
115/// `liveness_title` is the title recorded at enter time. `alive` and
116/// `current_title` are filled in by the supervisor from its `PROCS` singleton
117/// so the client does not need process introspection.
118#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
119pub struct ProjectSessionInfo {
120 pub pid: u32,
121 pub directory: PathBuf,
122 #[serde(skip_serializing_if = "Option::is_none", default)]
123 pub liveness_title: Option<String>,
124 pub alive: bool,
125 #[serde(skip_serializing_if = "Option::is_none", default)]
126 pub current_title: Option<String>,
127}
128
129#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, strum::Display, strum::EnumIs)]
130pub enum IpcResponse {
131 Ok,
132 /// Successful connect handshake, includes supervisor version for mismatch detection
133 ConnectOk {
134 version: String,
135 },
136 Yes,
137 No,
138 Error(String),
139 Notifications(Vec<(log::LevelFilter, String)>),
140 ActiveDaemons(Vec<Daemon>),
141 DisabledDaemons(Vec<DaemonId>),
142 DaemonAlreadyRunning,
143 DaemonStart {
144 daemon: Daemon,
145 },
146 DaemonFailed {
147 error: String,
148 },
149 /// Port conflict detected with detailed process information
150 PortConflict {
151 port: u16,
152 process: String,
153 pid: u32,
154 },
155 /// No available ports found after exhausting auto-bump attempts
156 NoAvailablePort {
157 start_port: u16,
158 attempts: u32,
159 },
160 DaemonReady {
161 daemon: Daemon,
162 },
163 DaemonFailedWithCode {
164 exit_code: Option<i32>,
165 /// Ports resolved by the failed attempt, so in-process retry hooks
166 /// observe the attempt's actual (post-bump) ports.
167 #[serde(default)]
168 resolved_ports: Vec<u16>,
169 },
170 /// Process was not running but had a PID record (unexpected exit)
171 DaemonWasNotRunning,
172 /// mDNS sync completed (or was a no-op if LAN mode is disabled)
173 MdnsSynced,
174 /// Settings reloaded from config files
175 ConfigReloaded,
176 /// URL of the running web UI, or `None` if the web UI is not running
177 WebUrl {
178 url: Option<String>,
179 },
180 /// Failed to kill the process (still running)
181 DaemonStopFailed {
182 error: String,
183 },
184 /// Daemon exists but is not running (no PID)
185 DaemonNotRunning,
186 DaemonNotFound,
187 /// Snapshot of all project sessions (response to `GetProjectSessions`).
188 ProjectSessions(Vec<ProjectSessionInfo>),
189 /// Number of daemon registrations removed by `CleanFiltered`.
190 Cleaned {
191 count: u64,
192 },
193}
194fn fs_name(name: &str) -> Result<Name<'_>> {
195 // Unix: use a filesystem path for the AF_UNIX socket.
196 #[cfg(unix)]
197 {
198 let path = env::IPC_SOCK_DIR.join(name).with_extension("sock");
199 let fs_name = path.to_fs_name::<GenericFilePath>().into_diagnostic()?;
200 Ok(fs_name)
201 }
202 // Windows: named pipes use a flat namespace (\\.\pipe\<name>) that
203 // cannot contain path separators. Derive a unique pipe name from the
204 // state directory to preserve test isolation when multiple supervisors
205 // run concurrently with different PITCHFORK_STATE_DIR values.
206 //
207 // Use a hash of the state directory path rather than character replacement
208 // to guarantee injectivity: `C:\a.b` and `C:\a\b` would both flatten to
209 // `C--a-b` with the old approach, causing pipe name collisions.
210 #[cfg(windows)]
211 {
212 let state_dir = env::PITCHFORK_STATE_DIR.to_string_lossy();
213 // FNV-1a hash: deterministic, stable across Rust versions.
214 // DefaultHasher's algorithm is not guaranteed stable, which would
215 // break IPC if the CLI and supervisor were ever compiled with
216 // different toolchains.
217 let mut hash: u64 = 0xcbf29ce484222325;
218 for byte in state_dir.bytes() {
219 hash ^= byte as u64;
220 hash = hash.wrapping_mul(0x100000001b3);
221 }
222 let pipe_name = format!("pitchfork-{hash:016x}-{name}");
223 Ok(pipe_name
224 .to_ns_name::<GenericNamespaced>()
225 .into_diagnostic()?)
226 }
227}
228
229fn serialize<T: serde::Serialize>(msg: &T) -> Result<Vec<u8>> {
230 if *env::IPC_JSON {
231 serde_json::to_vec(msg)
232 .into_diagnostic()
233 .wrap_err("failed to serialize IPC message as JSON")
234 } else {
235 rmp_serde::to_vec(msg)
236 .into_diagnostic()
237 .wrap_err("failed to serialize IPC message as MessagePack")
238 }
239}
240
241fn deserialize<T: serde::de::DeserializeOwned>(bytes: &[u8]) -> Result<T> {
242 let mut bytes = bytes.to_vec();
243 bytes.pop();
244 let preview = std::str::from_utf8(&bytes).unwrap_or("<binary>");
245 trace!("msg: {preview:?}");
246 if *env::IPC_JSON {
247 serde_json::from_slice(&bytes)
248 .into_diagnostic()
249 .wrap_err("failed to deserialize IPC JSON response")
250 } else {
251 rmp_serde::from_slice(&bytes)
252 .into_diagnostic()
253 .wrap_err("failed to deserialize IPC MessagePack response")
254 }
255}
256
257#[cfg(test)]
258mod tests {
259 use super::*;
260
261 #[test]
262 fn filtered_clean_ipc_round_trips() {
263 let request = IpcRequest::CleanFiltered {
264 namespaces: vec!["worktree".to_string()],
265 daemons: vec![DaemonId::new("worktree", "api")],
266 prune: true,
267 };
268 let mut bytes = serialize(&request).unwrap();
269 bytes.push(b'\n');
270 let decoded: IpcRequest = deserialize(&bytes).unwrap();
271 match decoded {
272 IpcRequest::CleanFiltered {
273 namespaces,
274 daemons,
275 prune,
276 } => {
277 assert_eq!(namespaces, ["worktree"]);
278 assert_eq!(daemons, [DaemonId::new("worktree", "api")]);
279 assert!(prune);
280 }
281 other => panic!("unexpected request: {other:?}"),
282 }
283
284 let mut bytes = serialize(&IpcResponse::Cleaned { count: 3 }).unwrap();
285 bytes.push(b'\n');
286 let decoded: IpcResponse = deserialize(&bytes).unwrap();
287 assert!(matches!(decoded, IpcResponse::Cleaned { count: 3 }));
288 }
289}