use crate::Result;
use crate::daemon::{Daemon, RunOptions};
use crate::daemon_id::DaemonId;
use crate::env;
use interprocess::local_socket::Name;
#[cfg(unix)]
use interprocess::local_socket::{GenericFilePath, ToFsName};
#[cfg(windows)]
use interprocess::local_socket::{GenericNamespaced, ToNsName};
use miette::{Context, IntoDiagnostic};
use std::path::PathBuf;
pub(crate) mod batch;
pub(crate) mod client;
pub(crate) mod server;
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, strum::Display, strum::EnumIs)]
#[allow(clippy::large_enum_variant)]
pub enum IpcRequest {
Connect,
ConnectV2 {
version: String,
},
Clean,
Stop {
id: DaemonId,
},
GetActiveDaemons,
GetDisabledDaemons,
Run(RunOptions),
Enable {
id: DaemonId,
},
Disable {
id: DaemonId,
},
UpdateShellDir {
shell_pid: u32,
dir: PathBuf,
},
GetNotifications,
SyncMdns,
ReloadConfig,
ProjectEnter {
pid: u32,
dir: PathBuf,
},
ProjectLeave {
pid: u32,
dir: PathBuf,
},
GetProjectSessions,
SinkOutputLine {
id: DaemonId,
token: u64,
fires_hook: bool,
line: String,
},
GetWebUrl,
CleanFiltered {
namespaces: Vec<String>,
daemons: Vec<DaemonId>,
prune: bool,
},
#[serde(skip)]
Invalid {
error: String,
},
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct ProjectSessionInfo {
pub pid: u32,
pub directory: PathBuf,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub liveness_title: Option<String>,
pub alive: bool,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub current_title: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, strum::Display, strum::EnumIs)]
pub enum IpcResponse {
Ok,
ConnectOk {
version: String,
},
Yes,
No,
Error(String),
Notifications(Vec<(log::LevelFilter, String)>),
ActiveDaemons(Vec<Daemon>),
DisabledDaemons(Vec<DaemonId>),
DaemonAlreadyRunning,
DaemonStart {
daemon: Daemon,
},
DaemonFailed {
error: String,
},
PortConflict {
port: u16,
process: String,
pid: u32,
},
NoAvailablePort {
start_port: u16,
attempts: u32,
},
DaemonReady {
daemon: Daemon,
},
DaemonFailedWithCode {
exit_code: Option<i32>,
#[serde(default)]
resolved_ports: Vec<u16>,
},
DaemonWasNotRunning,
MdnsSynced,
ConfigReloaded,
WebUrl {
url: Option<String>,
},
DaemonStopFailed {
error: String,
},
DaemonNotRunning,
DaemonNotFound,
ProjectSessions(Vec<ProjectSessionInfo>),
Cleaned {
count: u64,
},
}
fn fs_name(name: &str) -> Result<Name<'_>> {
#[cfg(unix)]
{
let path = env::IPC_SOCK_DIR.join(name).with_extension("sock");
let fs_name = path.to_fs_name::<GenericFilePath>().into_diagnostic()?;
Ok(fs_name)
}
#[cfg(windows)]
{
let state_dir = env::PITCHFORK_STATE_DIR.to_string_lossy();
let mut hash: u64 = 0xcbf29ce484222325;
for byte in state_dir.bytes() {
hash ^= byte as u64;
hash = hash.wrapping_mul(0x100000001b3);
}
let pipe_name = format!("pitchfork-{hash:016x}-{name}");
pipe_name
.to_ns_name::<GenericNamespaced>()
.into_diagnostic()
}
}
fn serialize<T: serde::Serialize>(msg: &T) -> Result<Vec<u8>> {
if *env::IPC_JSON {
serde_json::to_vec(msg)
.into_diagnostic()
.wrap_err("failed to serialize IPC message as JSON")
} else {
rmp_serde::to_vec(msg)
.into_diagnostic()
.wrap_err("failed to serialize IPC message as MessagePack")
}
}
fn deserialize<T: serde::de::DeserializeOwned>(bytes: &[u8]) -> Result<T> {
let mut bytes = bytes.to_vec();
bytes.pop();
let preview = std::str::from_utf8(&bytes).unwrap_or("<binary>");
trace!("msg: {preview:?}");
if *env::IPC_JSON {
serde_json::from_slice(&bytes)
.into_diagnostic()
.wrap_err("failed to deserialize IPC JSON response")
} else {
rmp_serde::from_slice(&bytes)
.into_diagnostic()
.wrap_err("failed to deserialize IPC MessagePack response")
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn filtered_clean_ipc_round_trips() {
let request = IpcRequest::CleanFiltered {
namespaces: vec!["worktree".to_string()],
daemons: vec![DaemonId::new("worktree", "api")],
prune: true,
};
let mut bytes = serialize(&request).unwrap();
bytes.push(b'\n');
let decoded: IpcRequest = deserialize(&bytes).unwrap();
match decoded {
IpcRequest::CleanFiltered {
namespaces,
daemons,
prune,
} => {
assert_eq!(namespaces, ["worktree"]);
assert_eq!(daemons, [DaemonId::new("worktree", "api")]);
assert!(prune);
}
other => panic!("unexpected request: {other:?}"),
}
let mut bytes = serialize(&IpcResponse::Cleaned { count: 3 }).unwrap();
bytes.push(b'\n');
let decoded: IpcResponse = deserialize(&bytes).unwrap();
assert!(matches!(decoded, IpcResponse::Cleaned { count: 3 }));
}
}