use crate::error::Result;
use crate::json_api::{self, JsonDoctorReport, JsonPath, JsonWorktree};
use crate::{config::Config, doctor, worktree};
use serde::Deserialize;
use serde_json::{json, Value};
use std::path::Path;
pub const PARSE_ERROR: i64 = -32700;
pub const INVALID_REQUEST: i64 = -32600;
pub const METHOD_NOT_FOUND: i64 = -32601;
pub const INVALID_PARAMS: i64 = -32602;
pub const INTERNAL_ERROR: i64 = -32603;
#[derive(Debug, Clone, Deserialize)]
pub struct RpcRequest {
#[serde(default)]
pub jsonrpc: String,
pub method: String,
#[serde(default)]
pub params: Value,
#[serde(default)]
pub id: Value,
}
pub fn success(id: &Value, result: Value) -> Value {
json!({ "jsonrpc": "2.0", "result": result, "id": id })
}
pub fn error(id: &Value, code: i64, message: &str) -> Value {
json!({ "jsonrpc": "2.0", "error": { "code": code, "message": message }, "id": id })
}
fn open_repo(workdir: &Path) -> Result<git2::Repository> {
worktree::discover_repo(Some(workdir))
}
fn run_list(workdir: &Path) -> Result<Vec<JsonWorktree>> {
let repo = open_repo(workdir)?;
json_api::worktrees(&repo)
}
fn run_path(workdir: &Path, pattern: &str) -> Result<JsonPath> {
let repo = open_repo(workdir)?;
let found = worktree::find_fuzzy(&repo, pattern)?;
Ok(JsonPath::from(&found))
}
fn run_doctor(workdir: &Path) -> Result<JsonDoctorReport> {
let repo = open_repo(workdir)?;
let repo_workdir = repo
.workdir()
.ok_or(crate::error::GwmError::NotInGitRepo)?
.to_path_buf();
let config = Config::load_for_repo(&repo_workdir).unwrap_or_default();
let global = crate::config::global_config_path();
let ctx = doctor::DoctorCtx {
repo_workdir: &repo_workdir,
repo: &repo,
config: &config,
global_config_path: global.as_deref(),
};
Ok(JsonDoctorReport::from(&doctor::run(&ctx)?))
}
pub fn dispatch(workdir: &Path, req: &RpcRequest) -> Value {
let id = &req.id;
match req.method.as_str() {
"list" => match run_list(workdir) {
Ok(list) => match serde_json::to_value(list) {
Ok(v) => success(id, v),
Err(e) => error(id, INTERNAL_ERROR, &e.to_string()),
},
Err(e) => error(id, INTERNAL_ERROR, &e.to_string()),
},
"doctor" => match run_doctor(workdir) {
Ok(report) => match serde_json::to_value(report) {
Ok(v) => success(id, v),
Err(e) => error(id, INTERNAL_ERROR, &e.to_string()),
},
Err(e) => error(id, INTERNAL_ERROR, &e.to_string()),
},
"path" => match req.params.get("pattern").and_then(|v| v.as_str()) {
None => error(id, INVALID_PARAMS, "method 'path' requires a string 'pattern' param"),
Some(pattern) => match run_path(workdir, pattern) {
Ok(p) => match serde_json::to_value(p) {
Ok(v) => success(id, v),
Err(e) => error(id, INTERNAL_ERROR, &e.to_string()),
},
Err(e) => error(id, INTERNAL_ERROR, &e.to_string()),
},
},
"subscribe" => error(
id,
INVALID_PARAMS,
"method 'subscribe' is only valid over a streaming socket connection",
),
other => error(id, METHOD_NOT_FOUND, &format!("unknown method '{other}'")),
}
}
pub fn handle_line(workdir: &Path, line: &str) -> Option<String> {
let value: Value = match serde_json::from_str(line) {
Ok(v) => v,
Err(e) => return Some(error(&Value::Null, PARSE_ERROR, &format!("parse error: {e}")).to_string()),
};
let req: RpcRequest = match serde_json::from_value(value.clone()) {
Ok(r) => r,
Err(e) => return Some(error(&Value::Null, INVALID_REQUEST, &format!("invalid request: {e}")).to_string()),
};
value.get("id")?;
Some(dispatch(workdir, &req).to_string())
}
pub fn worktrees_changed_notification(worktrees: &[JsonWorktree]) -> Value {
json!({
"jsonrpc": "2.0",
"method": "worktrees.changed",
"params": {
"schema_version": crate::contract::SCHEMA_VERSION,
"worktrees": worktrees,
},
})
}
pub fn worktrees_differ(old: &[JsonWorktree], new: &[JsonWorktree]) -> bool {
if old.len() != new.len() {
return true;
}
old.iter().zip(new).any(|(a, b)| {
a.name != b.name
|| a.id != b.id
|| a.path != b.path
|| a.branch != b.branch
|| a.head != b.head
|| a.is_main != b.is_main
|| a.is_locked != b.is_locked
|| a.is_prunable != b.is_prunable
|| a.status != b.status
|| a.issue != b.issue
|| a.pr != b.pr
|| a.agents != b.agents
})
}
pub fn next_subscription_push(
last: &Option<Vec<JsonWorktree>>,
latest: Result<Vec<JsonWorktree>>,
) -> Option<Vec<JsonWorktree>> {
let now = match latest {
Ok(now) => now,
Err(_) => return None,
};
match last {
None => Some(now), Some(prev) if worktrees_differ(prev, &now) => Some(now), Some(_) => None, }
}
pub const LIST_REQUEST: &str = r#"{"jsonrpc":"2.0","method":"list","id":1}"#;
pub const SUBSCRIBE_REQUEST: &str = r#"{"jsonrpc":"2.0","method":"subscribe","id":1}"#;
pub fn parse_list_result(line: &str) -> Result<Vec<JsonWorktree>> {
let v: Value = serde_json::from_str(line)
.map_err(|e| crate::error::GwmError::Other(format!("daemon: malformed list response: {e}")))?;
if let Some(err) = v.get("error") {
let msg = err.get("message").and_then(Value::as_str).unwrap_or("unknown error");
return Err(crate::error::GwmError::Other(format!("daemon list error: {msg}")));
}
let result = v
.get("result")
.ok_or_else(|| crate::error::GwmError::Other("daemon list response missing 'result'".into()))?;
serde_json::from_value(result.clone())
.map_err(|e| crate::error::GwmError::Other(format!("daemon: cannot decode worktree list: {e}")))
}
pub fn parse_worktrees_changed(line: &str) -> Result<Vec<JsonWorktree>> {
let v: Value = serde_json::from_str(line)
.map_err(|e| crate::error::GwmError::Other(format!("daemon: malformed notification: {e}")))?;
let arr = v
.get("params")
.and_then(|p| p.get("worktrees"))
.ok_or_else(|| crate::error::GwmError::Other("daemon notification missing 'params.worktrees'".into()))?;
serde_json::from_value(arr.clone())
.map_err(|e| crate::error::GwmError::Other(format!("daemon: cannot decode notification worktrees: {e}")))
}
#[cfg(all(unix, feature = "daemon"))]
pub mod client {
use super::*;
use crate::error::GwmError;
use std::io::{BufRead, BufReader, Write};
use std::os::unix::net::UnixStream;
use std::time::Duration;
const CLIENT_TIMEOUT: Duration = Duration::from_secs(5);
fn connect(socket: &Path, timeout: Option<Duration>) -> Result<UnixStream> {
let stream = UnixStream::connect(socket)
.map_err(|e| GwmError::Other(format!("daemon: cannot connect to {}: {e}", socket.display())))?;
stream
.set_read_timeout(timeout)
.map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
stream
.set_write_timeout(timeout)
.map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
Ok(stream)
}
pub fn list_once(socket: &Path) -> Result<Vec<JsonWorktree>> {
list_once_with_timeout(socket, Some(CLIENT_TIMEOUT))
}
#[doc(hidden)]
pub fn list_once_with_timeout(socket: &Path, timeout: Option<Duration>) -> Result<Vec<JsonWorktree>> {
let stream = connect(socket, timeout)?;
let mut writer = stream
.try_clone()
.map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
writeln!(writer, "{LIST_REQUEST}").map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
writer.flush().map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
let mut reader = BufReader::new(stream);
let mut line = String::new();
reader
.read_line(&mut line)
.map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
parse_list_result(line.trim())
}
pub fn subscribe(socket: &Path, mut on_snapshot: impl FnMut(&[JsonWorktree]) -> bool) -> Result<()> {
let stream = connect(socket, Some(CLIENT_TIMEOUT))?;
let mut writer = stream
.try_clone()
.map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
writeln!(writer, "{SUBSCRIBE_REQUEST}").map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
writer.flush().map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
let mut reader = BufReader::new(stream);
let mut delivered_any = false;
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) => break, Ok(_) => {}
Err(_) => break, }
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
let worktrees = parse_worktrees_changed(trimmed)?;
if !delivered_any {
let _ = reader.get_ref().set_read_timeout(None);
}
delivered_any = true;
if !on_snapshot(&worktrees) {
break;
}
}
if !delivered_any {
return Err(GwmError::Other(
"daemon: stream closed before the first snapshot".to_string(),
));
}
Ok(())
}
}
#[cfg(all(unix, feature = "daemon"))]
mod server {
use super::*;
use crate::error::GwmError;
use std::io::{BufRead, BufReader, ErrorKind, Read, Write};
use std::os::unix::fs::{DirBuilderExt, FileTypeExt, MetadataExt, PermissionsExt};
use std::os::unix::net::{UnixListener, UnixStream};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
struct ActiveGuard(Arc<AtomicUsize>);
impl ActiveGuard {
fn try_acquire(active: &Arc<AtomicUsize>, max: usize) -> Option<Self> {
if active.fetch_add(1, Ordering::SeqCst) + 1 > max {
active.fetch_sub(1, Ordering::SeqCst);
return None;
}
Some(ActiveGuard(Arc::clone(active)))
}
}
impl Drop for ActiveGuard {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::SeqCst);
}
}
const ACCEPT_TICK: Duration = Duration::from_millis(50);
pub struct ServeOptions {
pub socket: PathBuf,
pub repo_workdir: PathBuf,
pub poll_interval: Duration,
pub max_line_len: usize,
pub read_timeout: Option<Duration>,
pub max_connections: usize,
pub manage_socket_dir: bool,
}
impl ServeOptions {
pub const DEFAULT_MAX_LINE_LEN: usize = 64 * 1024;
pub const DEFAULT_READ_TIMEOUT: Duration = Duration::from_secs(30);
pub const DEFAULT_MAX_CONNECTIONS: usize = 128;
pub fn new(socket: PathBuf, repo_workdir: PathBuf, poll_interval: Duration) -> Self {
Self {
socket,
repo_workdir,
poll_interval,
max_line_len: Self::DEFAULT_MAX_LINE_LEN,
read_timeout: Some(Self::DEFAULT_READ_TIMEOUT),
max_connections: Self::DEFAULT_MAX_CONNECTIONS,
manage_socket_dir: false,
}
}
}
fn current_uid() -> u32 {
unsafe { libc::getuid() }
}
fn private_subdir_name() -> String {
format!("gwm-{}", current_uid())
}
fn is_private_dir(dir: &Path) -> bool {
match std::fs::symlink_metadata(dir) {
Ok(m) => m.file_type().is_dir() && m.uid() == current_uid() && m.mode() & 0o077 == 0,
Err(_) => false,
}
}
pub fn socket_in(base: &Path) -> PathBuf {
if is_private_dir(base) {
base.join("gwm.sock")
} else {
base.join(private_subdir_name()).join("gwm.sock")
}
}
pub fn socket_path() -> PathBuf {
if let Some(base) = std::env::var_os("XDG_RUNTIME_DIR").filter(|s| !s.is_empty()) {
return socket_in(&PathBuf::from(base));
}
if let Some(base) = std::env::var_os("TMPDIR").filter(|s| !s.is_empty()) {
return socket_in(&PathBuf::from(base));
}
socket_in(Path::new("/tmp"))
}
pub fn default_socket() -> (PathBuf, bool) {
let path = socket_path();
let managed =
path.parent().and_then(|d| d.file_name()).and_then(|n| n.to_str()) == Some(private_subdir_name().as_str());
(path, managed)
}
fn ensure_private_dir(dir: &Path) -> Result<()> {
let meta = match std::fs::symlink_metadata(dir) {
Ok(m) => m,
Err(_) => {
return std::fs::DirBuilder::new()
.mode(0o700)
.create(dir)
.map_err(|e| GwmError::Other(format!("daemon: failed to create private dir {}: {e}", dir.display())));
}
};
if !meta.file_type().is_dir() {
return Err(GwmError::Other(format!(
"daemon: refusing to use {}: exists and is not a directory",
dir.display()
)));
}
if meta.uid() != current_uid() {
return Err(GwmError::Other(format!(
"daemon: refusing to use {}: not owned by the current user",
dir.display()
)));
}
if meta.mode() & 0o077 != 0 {
std::fs::set_permissions(dir, std::fs::Permissions::from_mode(0o700))
.map_err(|e| GwmError::Other(format!("daemon: failed to restrict perms on {}: {e}", dir.display())))?;
}
Ok(())
}
fn clear_stale_socket(path: &Path) -> Result<()> {
let meta = match std::fs::symlink_metadata(path) {
Ok(m) => m,
Err(_) => return Ok(()), };
if !meta.file_type().is_socket() {
return Err(GwmError::Other(format!(
"daemon: refusing to use {}: exists and is not a unix socket",
path.display()
)));
}
if UnixStream::connect(path).is_ok() {
return Err(GwmError::Other(format!(
"daemon: socket {} is already in use by a live daemon",
path.display()
)));
}
let _ = std::fs::remove_file(path);
Ok(())
}
pub fn serve(opts: &ServeOptions, shutdown: Arc<AtomicBool>) -> Result<()> {
if opts.manage_socket_dir {
if let Some(parent) = opts.socket.parent() {
ensure_private_dir(parent)?;
}
}
clear_stale_socket(&opts.socket)?;
let listener = UnixListener::bind(&opts.socket)
.map_err(|e| GwmError::Other(format!("daemon: failed to bind {}: {e}", opts.socket.display())))?;
std::fs::set_permissions(&opts.socket, std::fs::Permissions::from_mode(0o600)).map_err(|e| {
let _ = std::fs::remove_file(&opts.socket);
GwmError::Other(format!(
"daemon: failed to restrict permissions on {}: {e}",
opts.socket.display()
))
})?;
listener
.set_nonblocking(true)
.map_err(|e| GwmError::Other(format!("daemon: set_nonblocking failed: {e}")))?;
let active = Arc::new(AtomicUsize::new(0));
eprintln!("gwm daemon listening on {}", opts.socket.display());
loop {
if shutdown.load(Ordering::Relaxed) {
break;
}
match listener.accept() {
Ok((stream, _addr)) => {
if let Err(e) = stream.set_nonblocking(false) {
eprintln!("daemon: failed to set connection blocking: {e}");
continue;
}
let Some(guard) = ActiveGuard::try_acquire(&active, opts.max_connections) else {
continue;
};
let workdir = opts.repo_workdir.clone();
let poll = opts.poll_interval;
let max_line_len = opts.max_line_len;
let read_timeout = opts.read_timeout;
let shutdown = Arc::clone(&shutdown);
std::thread::spawn(move || {
let _guard = guard;
handle_connection(stream, &workdir, poll, max_line_len, read_timeout, &shutdown);
});
}
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
std::thread::sleep(ACCEPT_TICK);
}
Err(e) => {
eprintln!("daemon: accept error: {e}");
std::thread::sleep(ACCEPT_TICK);
}
}
}
let _ = std::fs::remove_file(&opts.socket);
Ok(())
}
fn read_request_line(
stream: &UnixStream,
reader: &mut BufReader<UnixStream>,
max_len: usize,
deadline: Option<Instant>,
) -> Option<Vec<u8>> {
let mut buf = Vec::new();
loop {
if let Some(dl) = deadline {
match dl.checked_duration_since(Instant::now()) {
Some(rem) if !rem.is_zero() => {
let _ = stream.set_read_timeout(Some(rem));
}
_ => return None, }
}
let chunk = match reader.fill_buf() {
Ok(c) => c,
Err(e) if e.kind() == ErrorKind::Interrupted => continue,
Err(_) => return None, };
if chunk.is_empty() {
return None; }
if let Some(pos) = chunk.iter().position(|&b| b == b'\n') {
if buf.len() + pos > max_len {
return None; }
buf.extend_from_slice(&chunk[..pos]);
reader.consume(pos + 1);
return Some(buf);
}
if buf.len() + chunk.len() > max_len {
return None; }
let n = chunk.len();
buf.extend_from_slice(chunk);
reader.consume(n);
}
}
fn handle_connection(
stream: UnixStream,
workdir: &Path,
poll: Duration,
max_line_len: usize,
read_timeout: Option<Duration>,
shutdown: &AtomicBool,
) {
let _ = stream.set_read_timeout(read_timeout);
let read_half = match stream.try_clone() {
Ok(s) => s,
Err(_) => return,
};
let mut writer = stream;
let mut reader = BufReader::new(read_half);
loop {
let deadline = read_timeout.map(|t| Instant::now() + t);
let Some(bytes) = read_request_line(&writer, &mut reader, max_line_len, deadline) else {
break; };
let line = match std::str::from_utf8(&bytes) {
Ok(s) => s.trim(),
Err(_) => break, };
if line.is_empty() {
continue;
}
let is_subscribe = serde_json::from_str::<RpcRequest>(line)
.map(|r| r.method == "subscribe")
.unwrap_or(false);
if is_subscribe {
stream_subscription(&mut writer, workdir, poll, shutdown);
return;
}
if let Some(response) = handle_line(workdir, line) {
if writeln!(writer, "{response}").is_err() || writer.flush().is_err() {
break;
}
}
}
}
fn stream_subscription(stream: &mut UnixStream, workdir: &Path, poll: Duration, shutdown: &AtomicBool) {
if stream.set_read_timeout(Some(poll)).is_err() {
return;
}
let mut last: Option<Vec<JsonWorktree>> = None;
if let Some(snapshot) = next_subscription_push(&last, run_list(workdir)) {
if send_notification(stream, &snapshot).is_err() {
return;
}
last = Some(snapshot);
}
let mut buf = [0u8; 64];
loop {
if shutdown.load(Ordering::Relaxed) {
return;
}
match stream.read(&mut buf) {
Ok(0) => return,
Ok(_) => {} Err(e) if matches!(e.kind(), ErrorKind::WouldBlock | ErrorKind::TimedOut) => {}
Err(_) => return,
}
if let Some(snapshot) = next_subscription_push(&last, run_list(workdir)) {
if send_notification(stream, &snapshot).is_err() {
return;
}
last = Some(snapshot);
}
}
}
fn send_notification(writer: &mut UnixStream, worktrees: &[JsonWorktree]) -> std::io::Result<()> {
let note = worktrees_changed_notification(worktrees);
writeln!(writer, "{note}")?;
writer.flush()
}
}
#[cfg(all(unix, feature = "daemon"))]
pub use server::{default_socket, serve, socket_in, socket_path, ServeOptions};
pub fn pipe_user_fragment(raw: &str) -> String {
let mut out = String::with_capacity(raw.len());
for c in raw.chars() {
if c.is_ascii_alphanumeric() {
out.push(c);
} else {
let mut buf = [0u8; 4];
for b in c.encode_utf8(&mut buf).bytes() {
out.push_str(&format!("_{b:02X}"));
}
}
}
out
}
#[cfg(all(windows, feature = "daemon"))]
mod server_win {
use super::*;
use crate::error::GwmError;
use interprocess::os::windows::named_pipe::{pipe_mode, PipeListenerOptions, PipeMode, PipeStream};
use interprocess::os::windows::security_descriptor::SecurityDescriptor;
use interprocess::ConnectWaitMode;
use std::io::{BufRead, BufReader, ErrorKind, Write};
use std::os::windows::io::{AsHandle, AsRawHandle};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
type Conn = PipeStream<pipe_mode::Bytes, pipe_mode::Bytes>;
struct ActiveGuard(Arc<AtomicUsize>);
impl ActiveGuard {
fn try_acquire(active: &Arc<AtomicUsize>, max: usize) -> Option<Self> {
if active.fetch_add(1, Ordering::SeqCst) + 1 > max {
active.fetch_sub(1, Ordering::SeqCst);
return None;
}
Some(ActiveGuard(Arc::clone(active)))
}
}
impl Drop for ActiveGuard {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::SeqCst);
}
}
const ACCEPT_TICK: Duration = Duration::from_millis(50);
const NB_TICK: Duration = Duration::from_millis(15);
pub struct ServeOptions {
pub socket: PathBuf,
pub repo_workdir: PathBuf,
pub poll_interval: Duration,
pub max_line_len: usize,
pub read_timeout: Option<Duration>,
pub max_connections: usize,
pub manage_socket_dir: bool,
}
impl ServeOptions {
pub const DEFAULT_MAX_LINE_LEN: usize = 64 * 1024;
pub const DEFAULT_READ_TIMEOUT: Duration = Duration::from_secs(30);
pub const DEFAULT_MAX_CONNECTIONS: usize = 128;
pub fn new(socket: PathBuf, repo_workdir: PathBuf, poll_interval: Duration) -> Self {
Self {
socket,
repo_workdir,
poll_interval,
max_line_len: Self::DEFAULT_MAX_LINE_LEN,
read_timeout: Some(Self::DEFAULT_READ_TIMEOUT),
max_connections: Self::DEFAULT_MAX_CONNECTIONS,
manage_socket_dir: false,
}
}
}
pub fn socket_path() -> PathBuf {
let user = std::env::var("USERNAME").unwrap_or_else(|_| "default".to_string());
let identity = match std::env::var("USERDOMAIN") {
Ok(domain) if !domain.is_empty() => format!("{domain}\\{user}"),
_ => user,
};
PathBuf::from(format!("gwm-{}.sock", pipe_user_fragment(&identity)))
}
pub fn default_socket() -> (PathBuf, bool) {
(socket_path(), false)
}
fn owner_only_descriptor() -> Result<SecurityDescriptor> {
let sddl = widestring::U16CString::from_str("D:P(A;;GA;;;OW)")
.map_err(|e| GwmError::Other(format!("daemon: cannot encode the pipe SDDL: {e}")))?;
SecurityDescriptor::deserialize(&sddl)
.map_err(|e| GwmError::Other(format!("daemon: cannot build the pipe security descriptor: {e}")))
}
fn peek_available(conn: &Conn) -> Option<usize> {
let mut avail: u32 = 0;
let ok = unsafe {
windows_sys::Win32::System::Pipes::PeekNamedPipe(
conn.as_handle().as_raw_handle(),
std::ptr::null_mut(),
0,
std::ptr::null_mut(),
&mut avail,
std::ptr::null_mut(),
)
};
if ok == 0 {
None } else {
Some(avail as usize)
}
}
pub fn serve(opts: &ServeOptions, shutdown: Arc<AtomicBool>) -> Result<()> {
let path = widestring::U16CString::from_str(format!("\\\\.\\pipe\\{}", opts.socket.display()))
.map_err(|e| GwmError::Other(format!("daemon: invalid pipe name {}: {e}", opts.socket.display())))?;
let probe = Conn::connect_by_path_with_wait_mode(path.as_ucstr(), ConnectWaitMode::Timeout(Duration::from_secs(1)));
if probe.is_ok() {
return Err(GwmError::Other(format!(
"daemon: pipe {} is already in use by a live daemon",
opts.socket.display()
)));
}
let mut options = PipeListenerOptions::new();
options.path = std::borrow::Cow::Owned(path);
options.mode = PipeMode::Bytes;
options.security_descriptor = Some(owner_only_descriptor()?);
let listener = options.create_duplex::<pipe_mode::Bytes>().map_err(|e| {
GwmError::Other(format!(
"daemon: failed to bind pipe {} (a name that is already claimed is refused — first-instance guard): {e}",
opts.socket.display()
))
})?;
listener
.set_nonblocking(true)
.map_err(|e| GwmError::Other(format!("daemon: set_nonblocking failed: {e}")))?;
let active = Arc::new(AtomicUsize::new(0));
eprintln!("gwm daemon listening on \\\\.\\pipe\\{}", opts.socket.display());
loop {
if shutdown.load(Ordering::Relaxed) {
break;
}
match listener.accept() {
Ok(conn) => {
if let Err(e) = conn.set_nonblocking(false) {
eprintln!("daemon: failed to set connection blocking: {e}");
continue;
}
let Some(guard) = ActiveGuard::try_acquire(&active, opts.max_connections) else {
continue;
};
let workdir = opts.repo_workdir.clone();
let poll = opts.poll_interval;
let max_line_len = opts.max_line_len;
let read_timeout = opts.read_timeout;
let shutdown = Arc::clone(&shutdown);
std::thread::spawn(move || {
let _guard = guard;
handle_connection(&conn, &workdir, poll, max_line_len, read_timeout, &shutdown);
});
}
Err(ref e) if e.kind() == ErrorKind::WouldBlock => {
std::thread::sleep(ACCEPT_TICK);
}
Err(e) => {
eprintln!("daemon: accept error: {e}");
std::thread::sleep(ACCEPT_TICK);
}
}
}
Ok(())
}
fn read_request_line(
conn: &Conn,
reader: &mut BufReader<&Conn>,
max_len: usize,
deadline: Option<Instant>,
shutdown: &AtomicBool,
) -> Option<Vec<u8>> {
let mut buf = Vec::new();
loop {
if shutdown.load(Ordering::Relaxed) {
return None;
}
if let Some(dl) = deadline {
if Instant::now() >= dl {
return None; }
}
if reader.buffer().is_empty() {
match peek_available(conn) {
None => return None, Some(0) => {
std::thread::sleep(NB_TICK);
continue;
}
Some(_) => {} }
}
let chunk = match reader.fill_buf() {
Ok(c) => c,
Err(e) if e.kind() == ErrorKind::Interrupted => continue,
Err(_) => return None, };
if chunk.is_empty() {
return None; }
if let Some(pos) = chunk.iter().position(|&b| b == b'\n') {
if buf.len() + pos > max_len {
return None; }
buf.extend_from_slice(&chunk[..pos]);
reader.consume(pos + 1);
return Some(buf);
}
if buf.len() + chunk.len() > max_len {
return None; }
let n = chunk.len();
buf.extend_from_slice(chunk);
reader.consume(n);
}
}
fn write_frame(mut conn: &Conn, bytes: &[u8]) -> std::io::Result<()> {
conn.write_all(bytes)?;
conn.flush()
}
fn handle_connection(
conn: &Conn,
workdir: &Path,
poll: Duration,
max_line_len: usize,
read_timeout: Option<Duration>,
shutdown: &AtomicBool,
) {
let mut reader = BufReader::new(conn);
loop {
let deadline = read_timeout.map(|t| Instant::now() + t);
let Some(bytes) = read_request_line(conn, &mut reader, max_line_len, deadline, shutdown) else {
return; };
let line = match std::str::from_utf8(&bytes) {
Ok(s) => s.trim(),
Err(_) => return, };
if line.is_empty() {
continue;
}
let is_subscribe = serde_json::from_str::<RpcRequest>(line)
.map(|r| r.method == "subscribe")
.unwrap_or(false);
if is_subscribe {
stream_subscription(conn, &mut reader, workdir, poll, shutdown);
return;
}
if let Some(response) = handle_line(workdir, line) {
let mut frame = response.into_bytes();
frame.push(b'\n');
if write_frame(conn, &frame).is_err() {
return;
}
}
}
}
fn stream_subscription(
conn: &Conn,
reader: &mut BufReader<&Conn>,
workdir: &Path,
poll: Duration,
shutdown: &AtomicBool,
) {
let mut last: Option<Vec<JsonWorktree>> = None;
if let Some(snapshot) = next_subscription_push(&last, run_list(workdir)) {
if send_notification(conn, &snapshot).is_err() {
return;
}
last = Some(snapshot);
}
loop {
let tick_end = Instant::now() + poll;
loop {
if shutdown.load(Ordering::Relaxed) {
return;
}
if !reader.buffer().is_empty() {
return;
}
match peek_available(conn) {
None => return, Some(0) => {} Some(_) => return, }
let now = Instant::now();
if now >= tick_end {
break;
}
std::thread::sleep(NB_TICK.min(tick_end - now));
}
if let Some(snapshot) = next_subscription_push(&last, run_list(workdir)) {
if send_notification(conn, &snapshot).is_err() {
return;
}
last = Some(snapshot);
}
}
}
fn send_notification(conn: &Conn, worktrees: &[JsonWorktree]) -> std::io::Result<()> {
let mut frame = worktrees_changed_notification(worktrees).to_string().into_bytes();
frame.push(b'\n');
write_frame(conn, &frame)
}
}
#[cfg(all(windows, feature = "daemon"))]
pub use server_win::{default_socket, serve, socket_path, ServeOptions};
#[cfg(all(windows, feature = "daemon"))]
pub mod client {
use super::*;
use crate::error::GwmError;
use interprocess::os::windows::named_pipe::{pipe_mode, PipeStream};
use interprocess::ConnectWaitMode;
use std::io::{BufRead, BufReader, Write};
use std::os::windows::io::{AsHandle, AsRawHandle};
use std::sync::mpsc;
use std::time::Duration;
type Conn = PipeStream<pipe_mode::Bytes, pipe_mode::Bytes>;
const CLIENT_TIMEOUT: Duration = Duration::from_secs(5);
fn pipe_path(socket: &Path) -> Result<widestring::U16CString> {
widestring::U16CString::from_str(format!("\\\\.\\pipe\\{}", socket.display()))
.map_err(|e| GwmError::Other(format!("daemon: invalid pipe name {}: {e}", socket.display())))
}
fn verify_server_owner(conn: &Conn) -> Result<()> {
use windows_sys::Win32::Foundation::{CloseHandle, LocalFree};
use windows_sys::Win32::Security::Authorization::{GetSecurityInfo, SE_KERNEL_OBJECT};
use windows_sys::Win32::Security::{
CreateWellKnownSid, EqualSid, GetTokenInformation, TokenUser, WinBuiltinAdministratorsSid,
OWNER_SECURITY_INFORMATION, PSID, SECURITY_MAX_SID_SIZE, TOKEN_QUERY, TOKEN_USER,
};
use windows_sys::Win32::System::Threading::{GetCurrentProcess, OpenProcessToken};
let deny = |what: &str| GwmError::Other(format!("daemon: refusing untrusted pipe server ({what})"));
let mut owner: PSID = std::ptr::null_mut();
let mut descriptor = std::ptr::null_mut();
let status = unsafe {
GetSecurityInfo(
conn.as_handle().as_raw_handle(),
SE_KERNEL_OBJECT,
OWNER_SECURITY_INFORMATION,
&mut owner,
std::ptr::null_mut(),
std::ptr::null_mut(),
std::ptr::null_mut(),
&mut descriptor,
)
};
if status != 0 || owner.is_null() {
return Err(deny("cannot read the pipe owner"));
}
let result = (|| {
let mut token = std::ptr::null_mut();
if unsafe { OpenProcessToken(GetCurrentProcess(), TOKEN_QUERY, &mut token) } == 0 {
return Err(deny("cannot open the process token"));
}
let mut user_buf = [0u64; 32];
let mut len = 0u32;
let got = unsafe {
GetTokenInformation(
token,
TokenUser,
user_buf.as_mut_ptr().cast(),
(user_buf.len() * 8) as u32,
&mut len,
)
};
unsafe { CloseHandle(token) };
if got == 0 {
return Err(deny("cannot read the token user"));
}
let user_sid = unsafe { (*user_buf.as_ptr().cast::<TOKEN_USER>()).User.Sid };
if unsafe { EqualSid(owner, user_sid) } != 0 {
return Ok(());
}
let mut admin_buf = [0u64; (SECURITY_MAX_SID_SIZE as usize).div_ceil(8)];
let mut admin_len = (admin_buf.len() * 8) as u32;
let admin_ok = unsafe {
CreateWellKnownSid(
WinBuiltinAdministratorsSid,
std::ptr::null_mut(),
admin_buf.as_mut_ptr().cast(),
&mut admin_len,
)
};
if admin_ok != 0 && unsafe { EqualSid(owner, admin_buf.as_ptr().cast_mut().cast()) } != 0 {
return Ok(());
}
Err(deny("owned by another account"))
})();
unsafe { LocalFree(descriptor.cast()) };
result
}
fn connect(socket: &Path) -> Result<Conn> {
let path = pipe_path(socket)?;
let conn = Conn::connect_by_path_with_wait_mode(path.as_ucstr(), ConnectWaitMode::Timeout(CLIENT_TIMEOUT))
.map_err(|e| {
GwmError::Other(format!(
"daemon: cannot connect to \\\\.\\pipe\\{}: {e}",
socket.display()
))
})?;
verify_server_owner(&conn)?;
Ok(conn)
}
pub fn list_once(socket: &Path) -> Result<Vec<JsonWorktree>> {
list_once_with_timeout(socket, Some(CLIENT_TIMEOUT))
}
#[doc(hidden)]
pub fn list_once_with_timeout(socket: &Path, timeout: Option<Duration>) -> Result<Vec<JsonWorktree>> {
let socket = socket.to_path_buf();
let (tx, rx) = mpsc::channel();
std::thread::spawn(move || {
let _ = tx.send(round_trip(&socket));
});
match timeout {
Some(t) => rx
.recv_timeout(t)
.map_err(|_| GwmError::Other("daemon: timed out waiting for the response".to_string()))?,
None => rx
.recv()
.map_err(|_| GwmError::Other("daemon: client thread died".to_string()))?,
}
}
fn round_trip(socket: &Path) -> Result<Vec<JsonWorktree>> {
let conn = connect(socket)?;
let (recv, mut send) = conn.split();
writeln!(send, "{LIST_REQUEST}").map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
send.flush().map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
let mut reader = BufReader::new(recv);
let mut line = String::new();
reader
.read_line(&mut line)
.map_err(|e| GwmError::Other(format!("daemon: {e}")))?;
parse_list_result(line.trim())
}
pub fn subscribe(socket: &Path, mut on_snapshot: impl FnMut(&[JsonWorktree]) -> bool) -> Result<()> {
let socket = socket.to_path_buf();
let (tx, rx) = mpsc::channel::<Result<String>>();
let (half_tx, half_rx) = mpsc::channel();
std::thread::spawn(move || {
let conn = match connect(&socket) {
Ok(s) => s,
Err(e) => {
let _ = tx.send(Err(e));
return;
}
};
let (recv, mut send) = conn.split();
if let Err(e) = writeln!(send, "{SUBSCRIBE_REQUEST}").and_then(|()| send.flush()) {
let _ = tx.send(Err(GwmError::Other(format!("daemon: {e}"))));
return;
}
let _ = half_tx.send(send);
let mut reader = BufReader::new(recv);
loop {
let mut line = String::new();
match reader.read_line(&mut line) {
Ok(0) | Err(_) => break, Ok(_) => {
if tx.send(Ok(line)).is_err() {
break; }
}
}
}
});
let mut delivered_any = false;
let mut result = Ok(());
loop {
let msg = if delivered_any {
rx.recv().ok()
} else {
rx.recv_timeout(CLIENT_TIMEOUT).ok()
};
let Some(msg) = msg else { break };
let line = match msg {
Ok(line) => line,
Err(e) => {
result = Err(e);
break;
}
};
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
match parse_worktrees_changed(trimmed) {
Ok(worktrees) => {
delivered_any = true;
if !on_snapshot(&worktrees) {
break;
}
}
Err(e) => {
result = Err(e);
break;
}
}
}
if let Ok(mut send) = half_rx.try_recv() {
let _ = writeln!(send, "bye").and_then(|()| send.flush());
}
result?;
if !delivered_any {
return Err(GwmError::Other(
"daemon: stream closed before the first snapshot".to_string(),
));
}
Ok(())
}
}