| 1 | //! Authenticated local daemon transport: Unix socket and Windows named pipe. |
| 2 | //! |
| 3 | //! The desktop shell attaches to a long-lived `codewhale app-server --socket` |
| 4 | //! daemon over a local socket instead of a TCP port: local multi-client, |
| 5 | //! peer-credential auth, nothing to firewall (CORE-PROTOCOL spec §5). The |
| 6 | //! wire is *identical* to the `--stdio` transport — newline-delimited |
| 7 | //! JSON-RPC 2.0 driven by the same `crate::run_stdio_loop` — with exactly |
| 8 | //! one addition in front of it: a `daemon/attach` handshake that establishes |
| 9 | //! who this client is and whether it owns the daemon. |
| 10 | //! |
| 11 | //! # Endpoint resolution |
| 12 | //! |
| 13 | //! In precedence order (see [`resolve_socket_path`]): |
| 14 | //! |
| 15 | //! 1. an explicit path (`--socket-path`); |
| 16 | //! 2. `$CODEWHALE_HOME/run/daemon.sock` when `CODEWHALE_HOME` is set — an |
| 17 | //! explicit home is an isolation boundary, so its daemon must not collide |
| 18 | //! with the default one; |
| 19 | //! 3. `$XDG_RUNTIME_DIR/codewhale/daemon.sock`; |
| 20 | //! 4. macOS: `~/Library/Application Support/codewhale/daemon.sock`; |
| 21 | //! 5. `~/.codewhale/run/daemon.sock`. |
| 22 | //! |
| 23 | //! Windows derives a local named-pipe endpoint from the selected home and |
| 24 | //! current principal. It shares the same dispatcher and ownership handshake; |
| 25 | //! no supported platform falls back to TCP. Windows runtime proof is separate. |
| 26 | //! |
| 27 | //! # Ownership |
| 28 | //! |
| 29 | //! Hermes' claim model, server-side: a client attaches with `mode: "claim"` |
| 30 | //! (it spawned the daemon and will manage its lifetime) or `mode: "attach"` |
| 31 | //! (it found a healthy daemon and is a guest). Only the current owner may |
| 32 | //! `shutdown` the daemon; guests get `not_daemon_owner`. When the owner |
| 33 | //! disconnects the slot frees, so a relaunched shell can re-claim the daemon |
| 34 | //! it left running — sessions survive UI restarts because the daemon does. |
| 35 | |
| 36 | use std::path::{Path, PathBuf}; |
| 37 | |
| 38 | use codewhale_protocol::RuntimeOwnerReceipt; |
| 39 | use serde::{Deserialize, Serialize}; |
| 40 | |
| 41 | /// Basename of the daemon socket inside the Codewhale runtime directory. |
| 42 | pub const DAEMON_SOCKET_FILE_NAME: &str = "daemon.sock"; |
| 43 | |
| 44 | /// Historical reserved endpoint used in unsupported-platform diagnostics. |
| 45 | /// Actual Windows owner endpoints are scoped to the selected home/principal. |
| 46 | pub const WINDOWS_NAMED_PIPE: &str = r"\\.\pipe\codewhale-daemon"; |
| 47 | |
| 48 | /// JSON-RPC method a client must send first on a daemon-socket connection. |
| 49 | pub const ATTACH_METHOD: &str = "daemon/attach"; |
| 50 | |
| 51 | /// Longest socket path the kernel accepts (`sun_path` minus the NUL). |
| 52 | pub const MAX_SOCKET_PATH_BYTES: usize = if cfg!(any(target_os = "macos", target_os = "ios")) { |
| 53 | 103 |
| 54 | } else { |
| 55 | 107 |
| 56 | }; |
| 57 | |
| 58 | /// Typed failures of the daemon socket transport. |
| 59 | #[derive(Debug, thiserror::Error)] |
| 60 | pub enum DaemonSocketError { |
| 61 | /// The platform has no daemon socket implementation. Never a silent |
| 62 | /// fallback: the caller must pick another transport explicitly. |
| 63 | #[error( |
| 64 | "the daemon socket transport is not supported on {platform}; the reserved endpoint \ |
| 65 | there is the named pipe {planned_endpoint}, which is not implemented yet" |
| 66 | )] |
| 67 | UnsupportedPlatform { |
| 68 | platform: &'static str, |
| 69 | planned_endpoint: &'static str, |
| 70 | }, |
| 71 | /// No home directory (or runtime directory) to derive a default path from. |
| 72 | #[error( |
| 73 | "cannot resolve the Codewhale runtime directory for the daemon socket: no home directory" |
| 74 | )] |
| 75 | RuntimeDirUnavailable, |
| 76 | /// `CODEWHALE_HOME` is set but not a usable absolute path. |
| 77 | #[error("invalid CODEWHALE_HOME override: {0}")] |
| 78 | InvalidHomeOverride(String), |
| 79 | /// Unix socket paths are limited to roughly one hundred bytes. |
| 80 | #[error("daemon socket path {} is {len} bytes; this platform allows at most {max}", path.display())] |
| 81 | PathTooLong { |
| 82 | path: PathBuf, |
| 83 | len: usize, |
| 84 | max: usize, |
| 85 | }, |
| 86 | /// Something other than a socket already sits at the path. Refused so a |
| 87 | /// misconfigured path can never delete a user's file. |
| 88 | #[error("{} exists and is not a unix socket; refusing to remove it", path.display())] |
| 89 | NotASocket { path: PathBuf }, |
| 90 | /// A daemon answered on the socket: this one must not replace it. |
| 91 | #[error( |
| 92 | "a live listener already answers on {}; refusing to replace it (another codewhale daemon, or something else bound to this path)", |
| 93 | path.display() |
| 94 | )] |
| 95 | AlreadyRunning { path: PathBuf }, |
| 96 | /// The liveness probe neither connected nor was refused within the |
| 97 | /// budget. Refused rather than clobbered; remove the file by hand if the |
| 98 | /// old daemon is truly gone. |
| 99 | #[error("liveness probe of {} timed out; refusing to replace a socket that may be live", path.display())] |
| 100 | ProbeTimedOut { path: PathBuf }, |
| 101 | /// Filesystem or socket I/O failed. |
| 102 | #[error("{context} ({})", path.display())] |
| 103 | Io { |
| 104 | context: &'static str, |
| 105 | path: PathBuf, |
| 106 | #[source] |
| 107 | source: std::io::Error, |
| 108 | }, |
| 109 | /// The app-server state (config, state store, runtime) failed to build. |
| 110 | #[error("failed to build daemon state")] |
| 111 | State(#[source] anyhow::Error), |
| 112 | } |
| 113 | |
| 114 | /// How to start the daemon socket transport. |
| 115 | #[derive(Debug, Clone, Default)] |
| 116 | pub struct DaemonSocketOptions { |
| 117 | /// Explicit socket path; `None` resolves the platform default. |
| 118 | pub socket_path: Option<PathBuf>, |
| 119 | /// Explicit config file, like `app-server --config`. |
| 120 | pub config_path: Option<PathBuf>, |
| 121 | } |
| 122 | |
| 123 | /// Inputs to [`resolve_socket_path`], separated from the environment so the |
| 124 | /// precedence rules are a pure, testable function. |
| 125 | #[derive(Debug, Clone, Default)] |
| 126 | pub struct SocketPathInputs { |
| 127 | /// `--socket-path`. |
| 128 | pub explicit: Option<PathBuf>, |
| 129 | /// A valid explicit `CODEWHALE_HOME`. |
| 130 | pub codewhale_home_override: Option<PathBuf>, |
| 131 | /// `$XDG_RUNTIME_DIR`, when set and non-empty. |
| 132 | pub xdg_runtime_dir: Option<PathBuf>, |
| 133 | /// The user's home directory. |
| 134 | pub user_home: Option<PathBuf>, |
| 135 | /// Whether the macOS Application Support layout applies. |
| 136 | pub macos: bool, |
| 137 | } |
| 138 | |
| 139 | impl SocketPathInputs { |
| 140 | /// Capture the live environment. |
| 141 | pub fn from_environment(explicit: Option<PathBuf>) -> Result<Self, DaemonSocketError> { |
| 142 | let codewhale_home_override = codewhale_paths::codewhale_home_override() |
| 143 | .map_err(|err| DaemonSocketError::InvalidHomeOverride(err.to_string()))?; |
| 144 | let xdg_runtime_dir = std::env::var_os("XDG_RUNTIME_DIR") |
| 145 | .filter(|value| !value.is_empty()) |
| 146 | .map(PathBuf::from); |
| 147 | Ok(Self { |
| 148 | explicit, |
| 149 | codewhale_home_override, |
| 150 | xdg_runtime_dir, |
| 151 | user_home: codewhale_paths::user_home(), |
| 152 | macos: cfg!(target_os = "macos"), |
| 153 | }) |
| 154 | } |
| 155 | } |
| 156 | |
| 157 | /// Apply the precedence rules documented at the module level and enforce the |
| 158 | /// kernel's path-length limit. |
| 159 | pub fn resolve_socket_path(inputs: &SocketPathInputs) -> Result<PathBuf, DaemonSocketError> { |
| 160 | let path = if let Some(explicit) = inputs.explicit.clone() { |
| 161 | explicit |
| 162 | } else if let Some(home) = inputs.codewhale_home_override.clone() { |
| 163 | home.join("run").join(DAEMON_SOCKET_FILE_NAME) |
| 164 | } else if let Some(runtime_dir) = inputs.xdg_runtime_dir.clone() { |
| 165 | runtime_dir.join("codewhale").join(DAEMON_SOCKET_FILE_NAME) |
| 166 | } else { |
| 167 | let user_home = inputs |
| 168 | .user_home |
| 169 | .clone() |
| 170 | .ok_or(DaemonSocketError::RuntimeDirUnavailable)?; |
| 171 | if inputs.macos { |
| 172 | user_home |
| 173 | .join("Library") |
| 174 | .join("Application Support") |
| 175 | .join("codewhale") |
| 176 | .join(DAEMON_SOCKET_FILE_NAME) |
| 177 | } else { |
| 178 | user_home |
| 179 | .join(codewhale_paths::CODEWHALE_APP_DIR) |
| 180 | .join("run") |
| 181 | .join(DAEMON_SOCKET_FILE_NAME) |
| 182 | } |
| 183 | }; |
| 184 | let len = path.as_os_str().len(); |
| 185 | if len > MAX_SOCKET_PATH_BYTES { |
| 186 | return Err(DaemonSocketError::PathTooLong { |
| 187 | path, |
| 188 | len, |
| 189 | max: MAX_SOCKET_PATH_BYTES, |
| 190 | }); |
| 191 | } |
| 192 | Ok(path) |
| 193 | } |
| 194 | |
| 195 | /// The socket path this host would use with no explicit override. |
| 196 | #[cfg(unix)] |
| 197 | pub fn default_socket_path() -> Result<PathBuf, DaemonSocketError> { |
| 198 | resolve_socket_path(&SocketPathInputs::from_environment(None)?) |
| 199 | } |
| 200 | |
| 201 | /// The socket path this host would use with no explicit override. |
| 202 | #[cfg(not(any(unix, windows)))] |
| 203 | pub fn default_socket_path() -> Result<PathBuf, DaemonSocketError> { |
| 204 | Err(unsupported_platform()) |
| 205 | } |
| 206 | |
| 207 | /// Who is on the other end of a daemon-socket connection. |
| 208 | #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] |
| 209 | pub struct ClientIdentity { |
| 210 | /// Product name of the client, e.g. `codewhale-desktop`. |
| 211 | pub name: String, |
| 212 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 213 | pub version: Option<String>, |
| 214 | #[serde(default, skip_serializing_if = "Option::is_none")] |
| 215 | pub pid: Option<u32>, |
| 216 | } |
| 217 | |
| 218 | /// Ownership intent carried by `daemon/attach`. |
| 219 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)] |
| 220 | #[serde(rename_all = "snake_case")] |
| 221 | pub enum AttachMode { |
| 222 | /// A guest: use the daemon, never stop it. |
| 223 | #[default] |
| 224 | Attach, |
| 225 | /// The daemon's owner: may `shutdown`. Fails if a live owner exists. |
| 226 | Claim, |
| 227 | } |
| 228 | |
| 229 | /// Role granted by a successful attach. |
| 230 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] |
| 231 | #[serde(rename_all = "snake_case")] |
| 232 | pub enum AttachRole { |
| 233 | Owner, |
| 234 | Attached, |
| 235 | } |
| 236 | |
| 237 | /// `daemon/attach` params. |
| 238 | #[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)] |
| 239 | #[serde(rename_all = "snake_case")] |
| 240 | pub enum AttachFrontend { |
| 241 | #[default] |
| 242 | Control, |
| 243 | Acp, |
| 244 | Listener, |
| 245 | } |
| 246 | |
| 247 | #[derive(Debug, Clone, Deserialize)] |
| 248 | pub struct AttachParams { |
| 249 | #[serde(default)] |
| 250 | pub frontend: AttachFrontend, |
| 251 | #[serde(default)] |
| 252 | pub listener: Option<crate::RuntimeListenerSelection>, |
| 253 | #[serde(default)] |
| 254 | pub scope: Option<crate::RuntimeFrontendScope>, |
| 255 | #[serde(default)] |
| 256 | pub acp_model: Option<String>, |
| 257 | pub client: ClientIdentity, |
| 258 | #[serde(default)] |
| 259 | pub mode: AttachMode, |
| 260 | /// Bundle-skew guard: when set, the daemon refuses the attach unless its |
| 261 | /// own version string matches exactly. |
| 262 | #[serde(default)] |
| 263 | pub expect_daemon_version: Option<String>, |
| 264 | #[serde(default)] |
| 265 | pub expect_owner: Option<RuntimeOwnerReceipt>, |
| 266 | } |
| 267 | |
| 268 | /// The typed refusal every non-unix entry point returns. Unused in the unix |
| 269 | /// library build by construction; the tests pin its wording on every host. |
| 270 | #[cfg(any(test, not(any(unix, windows))))] |
| 271 | fn unsupported_platform() -> DaemonSocketError { |
| 272 | DaemonSocketError::UnsupportedPlatform { |
| 273 | platform: std::env::consts::OS, |
| 274 | planned_endpoint: WINDOWS_NAMED_PIPE, |
| 275 | } |
| 276 | } |
| 277 | |
| 278 | /// One bounded blocking boundary for held filesystem/process identity work. |
| 279 | /// The worker retains its permit and captured handles if the waiter cancels. |
| 280 | #[cfg(any(unix, windows))] |
| 281 | pub async fn owner_work<T, F>(work: F) -> anyhow::Result<T> |
| 282 | where |
| 283 | T: Send + 'static, |
| 284 | F: FnOnce() -> anyhow::Result<T> + Send + 'static, |
| 285 | { |
| 286 | static SLOTS: std::sync::OnceLock<std::sync::Arc<tokio::sync::Semaphore>> = |
| 287 | std::sync::OnceLock::new(); |
| 288 | let permit = SLOTS |
| 289 | .get_or_init(|| std::sync::Arc::new(tokio::sync::Semaphore::new(4))) |
| 290 | .clone() |
| 291 | .acquire_owned() |
| 292 | .await?; |
| 293 | tokio::task::spawn_blocking(move || { |
| 294 | let _permit = permit; |
| 295 | work() |
| 296 | }) |
| 297 | .await |
| 298 | .map_err(anyhow::Error::from)? |
| 299 | } |
| 300 | |
| 301 | use crate::{ |
| 302 | AppState, AppTransport, JsonRpcError, ParsedStdioLine, ShutdownAuthority, StdioLoopExit, |
| 303 | StdioLoopPolicy, dispatch_stdio_request_with_writer, jsonrpc_error, jsonrpc_result, |
| 304 | params_or_object, parse_params, parse_stdio_line, run_stdio_loop, write_stdio_line, |
| 305 | }; |
| 306 | use anyhow::Result; |
| 307 | use serde_json::{Value, json}; |
| 308 | use std::collections::HashMap; |
| 309 | use std::sync::{Arc, Mutex}; |
| 310 | use std::time::{Duration, Instant}; |
| 311 | use tokio::io::{AsyncBufRead, AsyncWrite, BufReader}; |
| 312 | use tokio::sync::watch; |
| 313 | use tokio::task::JoinSet; |
| 314 | |
| 315 | /// Facts about this daemon, reported in every attach reply. |
| 316 | #[derive(Debug)] |
| 317 | struct DaemonInfo { |
| 318 | pid: u32, |
| 319 | version: &'static str, |
| 320 | /// Independent effective principal captured from this process. |
| 321 | #[cfg(unix)] |
| 322 | uid: u32, |
| 323 | socket_path: PathBuf, |
| 324 | started_at: Instant, |
| 325 | } |
| 326 | |
| 327 | #[derive(Debug, Default)] |
| 328 | struct ConnectionRegistry { |
| 329 | next_id: u64, |
| 330 | connections: HashMap<u64, ClientIdentity>, |
| 331 | owner: Option<u64>, |
| 332 | } |
| 333 | |
| 334 | impl ConnectionRegistry { |
| 335 | fn register(&mut self, client: ClientIdentity) -> u64 { |
| 336 | self.next_id += 1; |
| 337 | let id = self.next_id; |
| 338 | self.connections.insert(id, client); |
| 339 | id |
| 340 | } |
| 341 | |
| 342 | /// Take the owner slot, or report who holds it. |
| 343 | fn claim(&mut self, id: u64) -> Result<(), ClientIdentity> { |
| 344 | if let Some(owner_id) = self.owner |
| 345 | && owner_id != id |
| 346 | && let Some(owner) = self.connections.get(&owner_id) |
| 347 | { |
| 348 | return Err(owner.clone()); |
| 349 | } |
| 350 | self.owner = Some(id); |
| 351 | Ok(()) |
| 352 | } |
| 353 | |
| 354 | fn owner(&self) -> Option<ClientIdentity> { |
| 355 | self.owner.and_then(|id| self.connections.get(&id)).cloned() |
| 356 | } |
| 357 | |
| 358 | fn remove(&mut self, id: u64) { |
| 359 | self.connections.remove(&id); |
| 360 | if self.owner == Some(id) { |
| 361 | self.owner = None; |
| 362 | } |
| 363 | } |
| 364 | } |
| 365 | |
| 366 | /// Releases the registry slot (and the owner claim) on drop, whichever |
| 367 | /// way the connection ends. |
| 368 | struct ConnectionGuard { |
| 369 | frontend: AttachFrontend, |
| 370 | listener: Option<crate::RuntimeListenerSelection>, |
| 371 | scope: Option<crate::RuntimeFrontendScope>, |
| 372 | acp_model: Option<String>, |
| 373 | registry: Arc<Mutex<ConnectionRegistry>>, |
| 374 | id: u64, |
| 375 | role: AttachRole, |
| 376 | } |
| 377 | |
| 378 | impl Drop for ConnectionGuard { |
| 379 | fn drop(&mut self) { |
| 380 | if let Ok(mut registry) = self.registry.lock() { |
| 381 | registry.remove(self.id); |
| 382 | } |
| 383 | } |
| 384 | } |
| 385 | |
| 386 | /// Asks a running [`DaemonSocket::serve`] to stop. |
| 387 | #[derive(Debug, Clone)] |
| 388 | pub struct DaemonShutdownHandle(pub(crate) Arc<watch::Sender<bool>>); |
| 389 | |
| 390 | impl DaemonShutdownHandle { |
| 391 | /// Idempotent; safe to call from any task or signal handler. |
| 392 | pub fn trigger(&self) { |
| 393 | self.0.send_replace(true); |
| 394 | } |
| 395 | } |
| 396 | |
| 397 | #[derive(Clone)] |
| 398 | pub(crate) struct ConnectionContext { |
| 399 | state: AppState, |
| 400 | registry: Arc<Mutex<ConnectionRegistry>>, |
| 401 | info: Arc<DaemonInfo>, |
| 402 | shutdown: DaemonShutdownHandle, |
| 403 | owner: Option<RuntimeOwnerReceipt>, |
| 404 | } |
| 405 | |
| 406 | /// The daemon's connection tasks. A `JoinSet` keeps every finished task |
| 407 | /// until it is joined, so a long-lived daemon used to retain one task per |
| 408 | /// past client (health polls, re-attaches) for its whole lifetime. |
| 409 | /// [`Self::spawn`] is the only way in and reaps finished tasks first, so |
| 410 | /// the set is bounded by the live connections plus those that ended since |
| 411 | /// the last accept. A panicked connection is logged, not fatal. |
| 412 | #[derive(Default)] |
| 413 | pub(crate) struct ConnectionTasks(JoinSet<()>); |
| 414 | |
| 415 | impl ConnectionTasks { |
| 416 | pub(crate) async fn shutdown(&mut self) { |
| 417 | self.0.shutdown().await; |
| 418 | } |
| 419 | pub(crate) fn at_capacity(&mut self) -> bool { |
| 420 | while let Some(joined) = self.0.try_join_next() { |
| 421 | if let Err(err) = joined |
| 422 | && err.is_panic() |
| 423 | { |
| 424 | tracing::warn!(error = %err, "daemon connection task panicked"); |
| 425 | } |
| 426 | } |
| 427 | self.0.len() >= 64 |
| 428 | } |
| 429 | pub(crate) fn spawn(&mut self, task: impl std::future::Future<Output = ()> + Send + 'static) { |
| 430 | if !self.at_capacity() { |
| 431 | self.0.spawn(task); |
| 432 | } else { |
| 433 | tracing::debug!("local control connection ceiling reached; refusing before admission"); |
| 434 | } |
| 435 | } |
| 436 | } |
| 437 | |
| 438 | /// Serve `healthz` and wait for `daemon/attach`; everything else is |
| 439 | /// refused with `attach_required` until the client attaches. |
| 440 | async fn handshake<R, W>( |
| 441 | context: &ConnectionContext, |
| 442 | lines: &mut crate::BoundedLines<R>, |
| 443 | writer: &mut W, |
| 444 | ) -> Result<Option<ConnectionGuard>> |
| 445 | where |
| 446 | R: AsyncBufRead + Unpin, |
| 447 | W: AsyncWrite + Unpin, |
| 448 | { |
| 449 | loop { |
| 450 | let Some(line) = lines.next_line().await? else { |
| 451 | return Ok(None); |
| 452 | }; |
| 453 | let request = match parse_stdio_line(&line) { |
| 454 | ParsedStdioLine::Blank => continue, |
| 455 | ParsedStdioLine::Rejected(response) => { |
| 456 | write_stdio_line(writer, &response).await?; |
| 457 | continue; |
| 458 | } |
| 459 | ParsedStdioLine::Request(request) => request, |
| 460 | }; |
| 461 | let id = request.id.clone(); |
| 462 | match request.method.as_str() { |
| 463 | "healthz" | "app/healthz" => { |
| 464 | let response = match dispatch_stdio_request_with_writer( |
| 465 | &context.state, |
| 466 | writer, |
| 467 | &request.method, |
| 468 | request.params, |
| 469 | AppTransport::Socket, |
| 470 | ) |
| 471 | .await |
| 472 | { |
| 473 | Ok(dispatch) => jsonrpc_result(id, dispatch.result), |
| 474 | Err(err) => jsonrpc_error(id, err), |
| 475 | }; |
| 476 | write_stdio_line(writer, &response).await?; |
| 477 | } |
| 478 | ATTACH_METHOD => match attach(context, request.params).await { |
| 479 | Ok((result, guard)) => { |
| 480 | write_stdio_line(writer, &jsonrpc_result(id, result)).await?; |
| 481 | return Ok(Some(guard)); |
| 482 | } |
| 483 | Err(err) => write_stdio_line(writer, &jsonrpc_error(id, err)).await?, |
| 484 | }, |
| 485 | other => { |
| 486 | write_stdio_line( |
| 487 | writer, |
| 488 | &jsonrpc_error(id, JsonRpcError::attach_required(other)), |
| 489 | ) |
| 490 | .await?; |
| 491 | } |
| 492 | } |
| 493 | } |
| 494 | } |
| 495 | |
| 496 | async fn attach( |
| 497 | context: &ConnectionContext, |
| 498 | params: Value, |
| 499 | ) -> Result<(Value, ConnectionGuard), JsonRpcError> { |
| 500 | let params: AttachParams = parse_params(params_or_object(params))?; |
| 501 | if params.frontend == AttachFrontend::Acp |
| 502 | && (context.state.owner_frontend.is_none() |
| 503 | || !context |
| 504 | .state |
| 505 | .captured_routing |
| 506 | .as_ref() |
| 507 | .is_some_and(|routing| routing.acp)) |
| 508 | { |
| 509 | return Err(JsonRpcError::invalid_params( |
| 510 | "held owner has no narrowed ACP projection; cross-profile attachment refused", |
| 511 | )); |
| 512 | } |
| 513 | if (params.frontend == AttachFrontend::Listener) != params.listener.is_some() { |
| 514 | return Err(JsonRpcError::invalid_params( |
| 515 | "listener selection does not match frontend", |
| 516 | )); |
| 517 | } |
| 518 | if let Some(listener) = params.listener.as_ref() { |
| 519 | listener.validate_bounds().map_err(|_| { |
| 520 | JsonRpcError::invalid_params("selected frontend input exceeds its bounds") |
| 521 | })?; |
| 522 | if context.state.owner_frontend.is_none() { |
| 523 | return Err(JsonRpcError::invalid_params( |
| 524 | "held owner has no selected listener projection", |
| 525 | )); |
| 526 | } |
| 527 | } |
| 528 | if params.scope.is_some() |
| 529 | && (params.frontend == AttachFrontend::Listener || context.state.owner_frontend.is_none()) |
| 530 | || params.acp_model.is_some() && params.frontend != AttachFrontend::Acp |
| 531 | { |
| 532 | return Err(JsonRpcError::invalid_params( |
| 533 | "selected scope does not match captured frontend", |
| 534 | )); |
| 535 | } |
| 536 | if let Some(scope) = params.scope.as_ref() { |
| 537 | scope.validate_bounds().map_err(|_| { |
| 538 | JsonRpcError::invalid_params("selected frontend scope exceeds its bounds") |
| 539 | })?; |
| 540 | } |
| 541 | if params |
| 542 | .acp_model |
| 543 | .as_ref() |
| 544 | .is_some_and(|model| model.trim().is_empty() || model.len() > 1024) |
| 545 | { |
| 546 | return Err(JsonRpcError::invalid_params("invalid selected ACP model")); |
| 547 | } |
| 548 | if let Some(owner) = context.owner.as_ref() { |
| 549 | if params.expect_owner.as_ref() != Some(owner) { |
| 550 | return Err(JsonRpcError::invalid_params( |
| 551 | "captured owner generation or store does not match", |
| 552 | )); |
| 553 | } |
| 554 | if params.mode == AttachMode::Claim { |
| 555 | return Err(JsonRpcError::not_daemon_owner()); |
| 556 | } |
| 557 | } |
| 558 | if params.client.name.trim().is_empty() { |
| 559 | return Err(JsonRpcError::invalid_params( |
| 560 | "client.name must not be empty", |
| 561 | )); |
| 562 | } |
| 563 | if let Some(expected) = params.expect_daemon_version.as_deref() |
| 564 | && expected != context.info.version |
| 565 | { |
| 566 | return Err(JsonRpcError::daemon_version_skew( |
| 567 | expected, |
| 568 | context.info.version, |
| 569 | )); |
| 570 | } |
| 571 | |
| 572 | if let Some(frontend) = context.state.owner_frontend.as_ref() { |
| 573 | let selection = match params.frontend { |
| 574 | AttachFrontend::Listener => Some(crate::RuntimeOwnerFrontendSelection::Listener( |
| 575 | params.listener.clone().expect("validated listener"), |
| 576 | )), |
| 577 | AttachFrontend::Acp => Some(crate::RuntimeOwnerFrontendSelection::Acp { |
| 578 | scope: params.scope.clone(), |
| 579 | model: params.acp_model.clone(), |
| 580 | }), |
| 581 | AttachFrontend::Control => params |
| 582 | .scope |
| 583 | .clone() |
| 584 | .map(crate::RuntimeOwnerFrontendSelection::Control), |
| 585 | }; |
| 586 | if let Some(selection) = selection { |
| 587 | frontend |
| 588 | .validate_selection(&selection) |
| 589 | .await |
| 590 | .map_err(|error| JsonRpcError::invalid_params(error.to_string()))?; |
| 591 | } |
| 592 | } |
| 593 | let mut registry = context |
| 594 | .registry |
| 595 | .lock() |
| 596 | .map_err(|_| JsonRpcError::internal("daemon connection registry poisoned"))?; |
| 597 | let id = registry.register(params.client.clone()); |
| 598 | let role = match params.mode { |
| 599 | AttachMode::Attach => AttachRole::Attached, |
| 600 | AttachMode::Claim => match registry.claim(id) { |
| 601 | Ok(()) => AttachRole::Owner, |
| 602 | Err(owner) => { |
| 603 | registry.remove(id); |
| 604 | let owner = serde_json::to_value(owner) |
| 605 | .map_err(|err| JsonRpcError::internal(err.to_string()))?; |
| 606 | return Err(JsonRpcError::daemon_already_claimed(&owner)); |
| 607 | } |
| 608 | }, |
| 609 | }; |
| 610 | let owner = registry.owner(); |
| 611 | let connections = registry.connections.len(); |
| 612 | drop(registry); |
| 613 | |
| 614 | let info = &context.info; |
| 615 | let result = json!({ |
| 616 | "attached": true, |
| 617 | "connection_id": id, |
| 618 | "role": role, |
| 619 | "transport": AppTransport::Socket.label(), |
| 620 | "daemon": { |
| 621 | "service": crate::legacy_deepseek_compat::SERVICE_NAME, |
| 622 | "pid": info.pid, |
| 623 | "version": info.version, |
| 624 | "socket_path": info.socket_path.display().to_string(), |
| 625 | "uptime_ms": u64::try_from(info.started_at.elapsed().as_millis()).unwrap_or(u64::MAX), |
| 626 | }, |
| 627 | "owner": owner, |
| 628 | "connections": connections, |
| 629 | "owner_receipt": context.owner, |
| 630 | "runtime_routing": context.state.captured_routing, |
| 631 | }); |
| 632 | Ok(( |
| 633 | result, |
| 634 | ConnectionGuard { |
| 635 | frontend: params.frontend, |
| 636 | listener: params.listener, |
| 637 | scope: params.scope, |
| 638 | acp_model: params.acp_model, |
| 639 | registry: Arc::clone(&context.registry), |
| 640 | id, |
| 641 | role, |
| 642 | }, |
| 643 | )) |
| 644 | } |
| 645 | |
| 646 | pub(crate) struct AuthorizedPeer { |
| 647 | pub(crate) pid: u32, |
| 648 | pub(crate) start: String, |
| 649 | #[cfg(windows)] |
| 650 | pub(crate) process: codewhale_config::windows_identity::WindowsPeerProcess, |
| 651 | } |
| 652 | |
| 653 | impl AuthorizedPeer { |
| 654 | pub(crate) fn check(&self) -> Result<()> { |
| 655 | #[cfg(unix)] |
| 656 | anyhow::ensure!( |
| 657 | codewhale_config::private_directory::unix_process_start(self.pid)? == self.start, |
| 658 | "local control peer generation changed" |
| 659 | ); |
| 660 | #[cfg(windows)] |
| 661 | { |
| 662 | self.process.check_current_user()?; |
| 663 | anyhow::ensure!( |
| 664 | self.process.pid() == self.pid, |
| 665 | "Windows kernel peer PID changed" |
| 666 | ); |
| 667 | anyhow::ensure!( |
| 668 | self.process.start() == self.start, |
| 669 | "local control peer generation changed" |
| 670 | ); |
| 671 | } |
| 672 | Ok(()) |
| 673 | } |
| 674 | } |
| 675 | |
| 676 | pub(crate) async fn connection_context( |
| 677 | state: AppState, |
| 678 | path: PathBuf, |
| 679 | shutdown: DaemonShutdownHandle, |
| 680 | owner: Option<RuntimeOwnerReceipt>, |
| 681 | ) -> Result<ConnectionContext> { |
| 682 | Ok(ConnectionContext { |
| 683 | state, |
| 684 | registry: Arc::new(Mutex::new(ConnectionRegistry::default())), |
| 685 | info: Arc::new(DaemonInfo { |
| 686 | pid: std::process::id(), |
| 687 | version: env!("CARGO_PKG_VERSION"), |
| 688 | #[cfg(unix)] |
| 689 | uid: codewhale_config::private_directory::PrivateDirectory::current_user_id(), |
| 690 | socket_path: path, |
| 691 | started_at: Instant::now(), |
| 692 | }), |
| 693 | shutdown, |
| 694 | owner, |
| 695 | }) |
| 696 | } |
| 697 | |
| 698 | pub(crate) async fn run_authorized_connection<R, W>( |
| 699 | context: ConnectionContext, |
| 700 | read: R, |
| 701 | mut writer: W, |
| 702 | peer: AuthorizedPeer, |
| 703 | ) where |
| 704 | R: tokio::io::AsyncRead + Unpin + Send + 'static, |
| 705 | W: AsyncWrite + Unpin + Send + 'static, |
| 706 | { |
| 707 | let peer = Arc::new(peer); |
| 708 | let check = peer.clone(); |
| 709 | if owner_work(move || check.check()).await.is_err() { |
| 710 | return; |
| 711 | } |
| 712 | let mut lines = crate::BoundedLines::new(BufReader::new(read)); |
| 713 | let guard = match tokio::time::timeout( |
| 714 | Duration::from_secs(5), |
| 715 | handshake(&context, &mut lines, &mut writer), |
| 716 | ) |
| 717 | .await |
| 718 | { |
| 719 | Err(_) => return, |
| 720 | Ok(result) => match result { |
| 721 | Ok(Some(guard)) => guard, |
| 722 | Ok(None) => return, |
| 723 | Err(error) => { |
| 724 | tracing::debug!(%error,"local control handshake ended"); |
| 725 | return; |
| 726 | } |
| 727 | }, |
| 728 | }; |
| 729 | let check = peer.clone(); |
| 730 | if owner_work(move || check.check()).await.is_err() { |
| 731 | return; |
| 732 | } |
| 733 | if guard.frontend != AttachFrontend::Control || guard.scope.is_some() { |
| 734 | let Some(frontend) = context.state.owner_frontend.as_ref() else { |
| 735 | return; |
| 736 | }; |
| 737 | let selection = match guard.frontend { |
| 738 | AttachFrontend::Listener => crate::RuntimeOwnerFrontendSelection::Listener( |
| 739 | guard.listener.clone().expect("admitted listener"), |
| 740 | ), |
| 741 | AttachFrontend::Acp => crate::RuntimeOwnerFrontendSelection::Acp { |
| 742 | scope: guard.scope.clone(), |
| 743 | model: guard.acp_model.clone(), |
| 744 | }, |
| 745 | AttachFrontend::Control => crate::RuntimeOwnerFrontendSelection::Control( |
| 746 | guard.scope.clone().expect("admitted scope"), |
| 747 | ), |
| 748 | }; |
| 749 | let _claim = (guard, peer); |
| 750 | if let Err(error) = frontend |
| 751 | .serve( |
| 752 | selection, |
| 753 | context.state.clone(), |
| 754 | Box::new(lines.into_inner()), |
| 755 | Box::new(writer), |
| 756 | ) |
| 757 | .await |
| 758 | { |
| 759 | tracing::debug!(%error,"captured frontend connection ended"); |
| 760 | } |
| 761 | return; |
| 762 | } |
| 763 | let policy = StdioLoopPolicy { |
| 764 | transport: AppTransport::Socket, |
| 765 | shutdown: match guard.role { |
| 766 | AttachRole::Owner => ShutdownAuthority::Granted, |
| 767 | AttachRole::Attached => ShutdownAuthority::Denied, |
| 768 | }, |
| 769 | }; |
| 770 | match run_stdio_loop(&context.state, lines, writer, policy, Some((guard, peer))).await { |
| 771 | Ok(StdioLoopExit::Shutdown) => context.shutdown.trigger(), |
| 772 | Ok(StdioLoopExit::InputClosed) => {} |
| 773 | Err(error) => tracing::debug!(%error,"local control connection ended"), |
| 774 | } |
| 775 | } |
| 776 | |
| 777 | #[cfg(unix)] |
| 778 | mod platform { |
| 779 | use std::fs::File; |
| 780 | use std::path::{Path, PathBuf}; |
| 781 | use std::sync::Arc; |
| 782 | use std::time::Duration; |
| 783 | |
| 784 | use anyhow::Result; |
| 785 | use codewhale_config::private_directory::{PrivateDirectory, PrivateSocketIdentity}; |
| 786 | use tokio::net::{UnixListener, UnixStream}; |
| 787 | use tokio::sync::watch; |
| 788 | |
| 789 | #[cfg(test)] |
| 790 | use super::{ClientIdentity, ConnectionRegistry}; |
| 791 | use super::{ |
| 792 | ConnectionContext, ConnectionTasks, DaemonShutdownHandle, DaemonSocketError, |
| 793 | DaemonSocketOptions, SocketPathInputs, resolve_socket_path, |
| 794 | }; |
| 795 | use crate::{AppState, AppTransport, build_state_off_runtime}; |
| 796 | |
| 797 | /// How long the stale-socket probe waits for a connect to resolve. |
| 798 | const PROBE_TIMEOUT: Duration = Duration::from_secs(1); |
| 799 | |
| 800 | /// Removes the socket file when the server stops, however it stops. |
| 801 | struct SocketFileGuard { |
| 802 | parent: Arc<PrivateDirectory>, |
| 803 | name: String, |
| 804 | identity: PrivateSocketIdentity, |
| 805 | receipt: Option<(String, File)>, |
| 806 | retired: bool, |
| 807 | } |
| 808 | |
| 809 | impl SocketFileGuard { |
| 810 | fn retirement(&mut self) -> Option<impl FnOnce() -> Result<()> + Send + 'static> { |
| 811 | if self.retired { |
| 812 | return None; |
| 813 | } |
| 814 | self.retired = true; |
| 815 | let parent = self.parent.clone(); |
| 816 | let name = self.name.clone(); |
| 817 | let identity = self.identity; |
| 818 | let receipt = self.receipt.take(); |
| 819 | Some(move || { |
| 820 | if let Some((name, file)) = receipt { |
| 821 | parent.retire_private_receipt(&name, &file)?; |
| 822 | } |
| 823 | parent.retire_socket(&name, identity)?; |
| 824 | Ok(()) |
| 825 | }) |
| 826 | } |
| 827 | |
| 828 | async fn retire(&mut self) -> Result<()> { |
| 829 | if let Some(retire) = self.retirement() { |
| 830 | tokio::spawn(async move { super::owner_work(retire).await }) |
| 831 | .await |
| 832 | .map_err(anyhow::Error::from)??; |
| 833 | } |
| 834 | Ok(()) |
| 835 | } |
| 836 | } |
| 837 | |
| 838 | impl Drop for SocketFileGuard { |
| 839 | fn drop(&mut self) { |
| 840 | let Some(retire) = self.retirement() else { |
| 841 | return; |
| 842 | }; |
| 843 | if let Ok(handle) = tokio::runtime::Handle::try_current() { |
| 844 | handle.spawn(async move { |
| 845 | if let Err(error) = super::owner_work(retire).await { |
| 846 | tracing::warn!(%error, "private endpoint retirement is uncertain; retaining it"); |
| 847 | } |
| 848 | }); |
| 849 | } else if let Err(error) = retire() { |
| 850 | tracing::warn!(%error, "private endpoint retirement is uncertain; retaining it"); |
| 851 | } |
| 852 | } |
| 853 | } |
| 854 | |
| 855 | /// A bound, not yet serving, daemon socket. |
| 856 | pub struct DaemonSocket { |
| 857 | listener: UnixListener, |
| 858 | path: PathBuf, |
| 859 | state: AppState, |
| 860 | shutdown: Arc<watch::Sender<bool>>, |
| 861 | socket_file: SocketFileGuard, |
| 862 | owner: Option<super::RuntimeOwnerReceipt>, |
| 863 | } |
| 864 | |
| 865 | impl DaemonSocket { |
| 866 | /// Where clients connect. |
| 867 | #[must_use] |
| 868 | pub fn local_path(&self) -> &Path { |
| 869 | &self.path |
| 870 | } |
| 871 | |
| 872 | /// A handle that stops [`Self::serve`] from outside (signals, tests). |
| 873 | #[must_use] |
| 874 | pub fn shutdown_handle(&self) -> DaemonShutdownHandle { |
| 875 | DaemonShutdownHandle(Arc::clone(&self.shutdown)) |
| 876 | } |
| 877 | |
| 878 | /// Accept clients until the owner sends `shutdown` or the handle is |
| 879 | /// triggered. Removes the socket file on the way out. |
| 880 | pub async fn serve(self) -> Result<(), DaemonSocketError> { |
| 881 | let Self { |
| 882 | listener, |
| 883 | path, |
| 884 | state, |
| 885 | shutdown, |
| 886 | socket_file, |
| 887 | owner, |
| 888 | } = self; |
| 889 | let mut socket_file = socket_file; |
| 890 | let context = super::connection_context( |
| 891 | state, |
| 892 | path.clone(), |
| 893 | DaemonShutdownHandle(Arc::clone(&shutdown)), |
| 894 | owner, |
| 895 | ) |
| 896 | .await |
| 897 | .map_err(DaemonSocketError::State)?; |
| 898 | let mut shutdown_rx = shutdown.subscribe(); |
| 899 | let mut connections = ConnectionTasks::default(); |
| 900 | |
| 901 | loop { |
| 902 | if *shutdown_rx.borrow() { |
| 903 | break; |
| 904 | } |
| 905 | tokio::select! { |
| 906 | accepted = listener.accept() => match accepted { |
| 907 | Ok((stream, _)) => { |
| 908 | connections.spawn(handle_connection(context.clone(), stream)); |
| 909 | } |
| 910 | Err(err) => { |
| 911 | tracing::warn!(error = %err, "daemon socket accept failed"); |
| 912 | tokio::time::sleep(Duration::from_millis(50)).await; |
| 913 | } |
| 914 | }, |
| 915 | changed = shutdown_rx.changed() => { |
| 916 | if changed.is_err() || *shutdown_rx.borrow() { |
| 917 | break; |
| 918 | } |
| 919 | } |
| 920 | } |
| 921 | } |
| 922 | |
| 923 | // The owner's `shutdown` reply was flushed before its loop |
| 924 | // returned, so aborting what is left loses nothing a client |
| 925 | // still needs. |
| 926 | connections.shutdown().await; |
| 927 | drop(listener); |
| 928 | socket_file |
| 929 | .retire() |
| 930 | .await |
| 931 | .map_err(DaemonSocketError::State)?; |
| 932 | Ok(()) |
| 933 | } |
| 934 | } |
| 935 | |
| 936 | /// Resolve the path, clear a stale socket, bind with `0600`, and build |
| 937 | /// the shared app state. Does not accept anything until |
| 938 | /// [`DaemonSocket::serve`]. |
| 939 | pub async fn bind_daemon_socket( |
| 940 | options: DaemonSocketOptions, |
| 941 | ) -> Result<DaemonSocket, DaemonSocketError> { |
| 942 | bind_with_owner(options, None).await |
| 943 | } |
| 944 | |
| 945 | pub(crate) async fn bind_captured_owner( |
| 946 | state: AppState, |
| 947 | owner: super::RuntimeOwnerReceipt, |
| 948 | ) -> Result<DaemonSocket, DaemonSocketError> { |
| 949 | let options = DaemonSocketOptions { |
| 950 | socket_path: Some(owner.socket_path.clone()), |
| 951 | config_path: None, |
| 952 | }; |
| 953 | bind_with_owner(options, Some((state, owner))).await |
| 954 | } |
| 955 | |
| 956 | async fn bind_with_owner( |
| 957 | options: DaemonSocketOptions, |
| 958 | captured: Option<(AppState, super::RuntimeOwnerReceipt)>, |
| 959 | ) -> Result<DaemonSocket, DaemonSocketError> { |
| 960 | let path = resolve_socket_path(&SocketPathInputs::from_environment(options.socket_path)?)?; |
| 961 | let parent_path = path |
| 962 | .parent() |
| 963 | .filter(|parent| !parent.as_os_str().is_empty()) |
| 964 | .ok_or(DaemonSocketError::RuntimeDirUnavailable)? |
| 965 | .to_path_buf(); |
| 966 | let parent = Arc::new( |
| 967 | super::owner_work(move || PrivateDirectory::admit(&parent_path)) |
| 968 | .await |
| 969 | .map_err(DaemonSocketError::State)?, |
| 970 | ); |
| 971 | let name = path |
| 972 | .file_name() |
| 973 | .and_then(|name| name.to_str()) |
| 974 | .ok_or(DaemonSocketError::RuntimeDirUnavailable)? |
| 975 | .to_string(); |
| 976 | clear_stale_socket(&path, &parent, &name).await?; |
| 977 | let check_parent = parent.clone(); |
| 978 | if !super::owner_work(move || check_parent.is_at_selected_path()) |
| 979 | .await |
| 980 | .map_err(DaemonSocketError::State)? |
| 981 | { |
| 982 | return Err(DaemonSocketError::State(anyhow::anyhow!( |
| 983 | "private endpoint parent changed" |
| 984 | ))); |
| 985 | } |
| 986 | |
| 987 | let listener = UnixListener::bind(&path).map_err(|source| DaemonSocketError::Io { |
| 988 | context: "failed to bind the daemon socket", |
| 989 | path: path.clone(), |
| 990 | source, |
| 991 | })?; |
| 992 | let inspect_parent = parent.clone(); |
| 993 | let inspect_name = name.clone(); |
| 994 | let identity = super::owner_work(move || inspect_parent.socket_identity(&inspect_name)) |
| 995 | .await |
| 996 | .map_err(DaemonSocketError::State)? |
| 997 | .ok_or_else(|| { |
| 998 | DaemonSocketError::State(anyhow::anyhow!("bound socket identity is unavailable")) |
| 999 | })?; |
| 1000 | let socket_file = SocketFileGuard { |
| 1001 | parent: parent.clone(), |
| 1002 | name, |
| 1003 | identity, |
| 1004 | receipt: None, |
| 1005 | retired: false, |
| 1006 | }; |
| 1007 | let protect_parent = parent.clone(); |
| 1008 | let protect_name = socket_file.name.clone(); |
| 1009 | super::owner_work(move || { |
| 1010 | anyhow::ensure!( |
| 1011 | protect_parent.is_at_selected_path()?, |
| 1012 | "private endpoint parent changed during bind" |
| 1013 | ); |
| 1014 | protect_parent.protect_socket(&protect_name, identity) |
| 1015 | }) |
| 1016 | .await |
| 1017 | .map_err(DaemonSocketError::State)?; |
| 1018 | |
| 1019 | let (state, owner) = match captured { |
| 1020 | Some((state, owner)) => (state, Some(owner)), |
| 1021 | None => ( |
| 1022 | build_state_off_runtime(options.config_path, None, AppTransport::Socket) |
| 1023 | .await |
| 1024 | .map_err(DaemonSocketError::State)?, |
| 1025 | None, |
| 1026 | ), |
| 1027 | }; |
| 1028 | let mut socket_file = socket_file; |
| 1029 | if let Some(owner) = owner.as_ref() { |
| 1030 | let receipt_name = format!("{}.owner.json", socket_file.name); |
| 1031 | let bytes = serde_json::to_vec(owner) |
| 1032 | .map_err(|error| DaemonSocketError::State(error.into()))?; |
| 1033 | if bytes.len() > 16384 { |
| 1034 | return Err(DaemonSocketError::State(anyhow::anyhow!( |
| 1035 | "owner receipt exceeds private publication limit" |
| 1036 | ))); |
| 1037 | } |
| 1038 | let parent = parent.clone(); |
| 1039 | let name = receipt_name.clone(); |
| 1040 | let expected = owner.clone(); |
| 1041 | let guard = socket_file; |
| 1042 | socket_file = super::owner_work(move || { |
| 1043 | let mut guard = guard; |
| 1044 | anyhow::ensure!( |
| 1045 | parent.is_at_selected_path()?, |
| 1046 | "private owner parent changed before publication" |
| 1047 | ); |
| 1048 | if let Some((bytes, old)) = parent.read_private_receipt(&name, 16384)? { |
| 1049 | let previous: super::RuntimeOwnerReceipt = serde_json::from_slice(&bytes)?; |
| 1050 | anyhow::ensure!( |
| 1051 | previous.version == expected.version |
| 1052 | && previous.data_dir == expected.data_dir |
| 1053 | && previous.execution_scope == expected.execution_scope |
| 1054 | && previous.socket_path == expected.socket_path, |
| 1055 | "stale owner receipt belongs to a different selected store" |
| 1056 | ); |
| 1057 | anyhow::ensure!( |
| 1058 | parent.retire_private_receipt(&name, &old)?, |
| 1059 | "stale owner receipt changed" |
| 1060 | ); |
| 1061 | } |
| 1062 | parent.write_owned_file(&name, &bytes, false)?; |
| 1063 | let (_, file) = parent |
| 1064 | .read_private_receipt(&name, 16384)? |
| 1065 | .ok_or_else(|| anyhow::anyhow!("published owner receipt unavailable"))?; |
| 1066 | guard.receipt = Some((name, file)); |
| 1067 | anyhow::ensure!( |
| 1068 | parent.is_at_selected_path()?, |
| 1069 | "private owner parent changed during publication; retaining uncertain receipt" |
| 1070 | ); |
| 1071 | Ok(guard) |
| 1072 | }) |
| 1073 | .await |
| 1074 | .map_err(DaemonSocketError::State)?; |
| 1075 | } |
| 1076 | let (shutdown, _) = watch::channel(false); |
| 1077 | Ok(DaemonSocket { |
| 1078 | listener, |
| 1079 | path, |
| 1080 | state, |
| 1081 | shutdown: Arc::new(shutdown), |
| 1082 | socket_file, |
| 1083 | owner, |
| 1084 | }) |
| 1085 | } |
| 1086 | |
| 1087 | /// Only a definite refusal permits exact captured endpoint retirement. |
| 1088 | async fn clear_stale_socket( |
| 1089 | path: &Path, |
| 1090 | parent: &Arc<PrivateDirectory>, |
| 1091 | name: &str, |
| 1092 | ) -> Result<(), DaemonSocketError> { |
| 1093 | let inspect_parent = parent.clone(); |
| 1094 | let inspect_name = name.to_string(); |
| 1095 | let Some(identity) = |
| 1096 | super::owner_work(move || inspect_parent.socket_identity(&inspect_name)) |
| 1097 | .await |
| 1098 | .map_err(|error| { |
| 1099 | if error |
| 1100 | .downcast_ref::<std::io::Error>() |
| 1101 | .is_some_and(|error| error.kind() == std::io::ErrorKind::InvalidInput) |
| 1102 | { |
| 1103 | DaemonSocketError::NotASocket { |
| 1104 | path: path.to_path_buf(), |
| 1105 | } |
| 1106 | } else { |
| 1107 | DaemonSocketError::State(error) |
| 1108 | } |
| 1109 | })? |
| 1110 | else { |
| 1111 | return Ok(()); |
| 1112 | }; |
| 1113 | match tokio::time::timeout(PROBE_TIMEOUT, UnixStream::connect(path)).await { |
| 1114 | Ok(Ok(_live)) => Err(DaemonSocketError::AlreadyRunning { |
| 1115 | path: path.to_path_buf(), |
| 1116 | }), |
| 1117 | Ok(Err(error)) if error.kind() == std::io::ErrorKind::ConnectionRefused => { |
| 1118 | let retire_parent = parent.clone(); |
| 1119 | let retire_name = name.to_string(); |
| 1120 | if super::owner_work(move || retire_parent.retire_socket(&retire_name, identity)) |
| 1121 | .await |
| 1122 | .map_err(DaemonSocketError::State)? |
| 1123 | { |
| 1124 | Ok(()) |
| 1125 | } else { |
| 1126 | Err(DaemonSocketError::State(anyhow::anyhow!( |
| 1127 | "stale endpoint changed; refusing replacement" |
| 1128 | ))) |
| 1129 | } |
| 1130 | } |
| 1131 | Ok(Err(source)) => Err(DaemonSocketError::Io { |
| 1132 | context: "uncertain daemon endpoint probe; refusing replacement", |
| 1133 | path: path.to_path_buf(), |
| 1134 | source, |
| 1135 | }), |
| 1136 | Err(_) => Err(DaemonSocketError::ProbeTimedOut { |
| 1137 | path: path.to_path_buf(), |
| 1138 | }), |
| 1139 | } |
| 1140 | } |
| 1141 | |
| 1142 | async fn handle_connection(context: ConnectionContext, stream: UnixStream) { |
| 1143 | let credential = match stream.peer_cred() { |
| 1144 | Ok(peer) if peer.uid() == context.info.uid => peer, |
| 1145 | _ => return, |
| 1146 | }; |
| 1147 | let Some(pid) = credential.pid().and_then(|pid| u32::try_from(pid).ok()) else { |
| 1148 | return; |
| 1149 | }; |
| 1150 | let start = match super::owner_work(move || { |
| 1151 | codewhale_config::private_directory::unix_process_start(pid) |
| 1152 | }) |
| 1153 | .await |
| 1154 | { |
| 1155 | Ok(start) => start, |
| 1156 | Err(_) => return, |
| 1157 | }; |
| 1158 | let (read, write) = stream.into_split(); |
| 1159 | super::run_authorized_connection( |
| 1160 | context, |
| 1161 | read, |
| 1162 | write, |
| 1163 | super::AuthorizedPeer { pid, start }, |
| 1164 | ) |
| 1165 | .await; |
| 1166 | } |
| 1167 | |
| 1168 | #[cfg(test)] |
| 1169 | mod tests { |
| 1170 | use super::*; |
| 1171 | |
| 1172 | fn client(name: &str) -> ClientIdentity { |
| 1173 | ClientIdentity { |
| 1174 | name: name.to_string(), |
| 1175 | version: None, |
| 1176 | pid: None, |
| 1177 | } |
| 1178 | } |
| 1179 | |
| 1180 | #[tokio::test] |
| 1181 | async fn finished_connection_tasks_are_reaped() { |
| 1182 | let mut connections = ConnectionTasks::default(); |
| 1183 | let mut handles = Vec::new(); |
| 1184 | for _ in 0..3 { |
| 1185 | handles.push(connections.0.spawn(async {})); |
| 1186 | } |
| 1187 | handles.push( |
| 1188 | connections |
| 1189 | .0 |
| 1190 | .spawn(async { panic!("fixture connection panic") }), |
| 1191 | ); |
| 1192 | tokio::time::timeout(Duration::from_secs(5), async { |
| 1193 | while !handles.iter().all(tokio::task::AbortHandle::is_finished) { |
| 1194 | tokio::task::yield_now().await; |
| 1195 | } |
| 1196 | }) |
| 1197 | .await |
| 1198 | .expect("fixture connections finish"); |
| 1199 | |
| 1200 | // Accepting the next client reaps every finished connection, |
| 1201 | // including the panicked one, and keeps only the new live one. |
| 1202 | connections.spawn(std::future::pending::<()>()); |
| 1203 | assert_eq!(connections.0.len(), 1, "only the live connection is kept"); |
| 1204 | connections.0.abort_all(); |
| 1205 | } |
| 1206 | |
| 1207 | #[tokio::test] |
| 1208 | async fn cancelled_owner_waiter_retains_worker_capacity_until_real_completion() { |
| 1209 | let (started, mut observed) = tokio::sync::mpsc::unbounded_channel(); |
| 1210 | let mut releases = Vec::new(); |
| 1211 | let mut waiters = Vec::new(); |
| 1212 | for _ in 0..4 { |
| 1213 | let (release, released) = std::sync::mpsc::channel(); |
| 1214 | releases.push(release); |
| 1215 | let started = started.clone(); |
| 1216 | waiters.push(tokio::spawn(super::super::owner_work(move || { |
| 1217 | started.send(()).unwrap(); |
| 1218 | let _ = released.recv(); |
| 1219 | Ok(()) |
| 1220 | }))); |
| 1221 | } |
| 1222 | for _ in 0..4 { |
| 1223 | tokio::time::timeout(Duration::from_secs(5), observed.recv()) |
| 1224 | .await |
| 1225 | .unwrap() |
| 1226 | .unwrap(); |
| 1227 | } |
| 1228 | waiters[0].abort(); |
| 1229 | let _ = (&mut waiters[0]).await; |
| 1230 | let (next_started, mut next_observed) = tokio::sync::mpsc::unbounded_channel(); |
| 1231 | let next = tokio::spawn(super::super::owner_work(move || { |
| 1232 | next_started.send(()).unwrap(); |
| 1233 | Ok(()) |
| 1234 | })); |
| 1235 | tokio::task::yield_now().await; |
| 1236 | assert!(matches!( |
| 1237 | next_observed.try_recv(), |
| 1238 | Err(tokio::sync::mpsc::error::TryRecvError::Empty) |
| 1239 | )); |
| 1240 | releases.remove(0).send(()).unwrap(); |
| 1241 | tokio::time::timeout(Duration::from_secs(5), next_observed.recv()) |
| 1242 | .await |
| 1243 | .unwrap() |
| 1244 | .unwrap(); |
| 1245 | for release in releases { |
| 1246 | release.send(()).unwrap(); |
| 1247 | } |
| 1248 | for waiter in waiters.into_iter().skip(1) { |
| 1249 | waiter.await.unwrap().unwrap(); |
| 1250 | } |
| 1251 | next.await.unwrap().unwrap(); |
| 1252 | } |
| 1253 | |
| 1254 | #[test] |
| 1255 | fn registry_claim_is_exclusive_until_the_owner_leaves() { |
| 1256 | let mut registry = ConnectionRegistry::default(); |
| 1257 | let first = registry.register(client("desktop-a")); |
| 1258 | let second = registry.register(client("desktop-b")); |
| 1259 | |
| 1260 | assert!(registry.claim(first).is_ok()); |
| 1261 | assert_eq!(registry.claim(second), Err(client("desktop-a"))); |
| 1262 | assert!( |
| 1263 | registry.claim(first).is_ok(), |
| 1264 | "re-claim by the owner is idempotent" |
| 1265 | ); |
| 1266 | |
| 1267 | registry.remove(first); |
| 1268 | assert_eq!(registry.owner(), None); |
| 1269 | assert!(registry.claim(second).is_ok()); |
| 1270 | assert_eq!(registry.owner(), Some(client("desktop-b"))); |
| 1271 | } |
| 1272 | |
| 1273 | #[test] |
| 1274 | fn removing_a_guest_keeps_the_owner() { |
| 1275 | let mut registry = ConnectionRegistry::default(); |
| 1276 | let owner = registry.register(client("owner")); |
| 1277 | let guest = registry.register(client("guest")); |
| 1278 | registry.claim(owner).expect("claim"); |
| 1279 | registry.remove(guest); |
| 1280 | assert_eq!(registry.owner(), Some(client("owner"))); |
| 1281 | assert_eq!(registry.connections.len(), 1); |
| 1282 | } |
| 1283 | } |
| 1284 | } |
| 1285 | |
| 1286 | #[cfg(not(any(unix, windows)))] |
| 1287 | mod platform { |
| 1288 | use std::path::Path; |
| 1289 | |
| 1290 | use super::{ |
| 1291 | DaemonShutdownHandle, DaemonSocketError, DaemonSocketOptions, unsupported_platform, |
| 1292 | }; |
| 1293 | |
| 1294 | /// Placeholder until the Windows named pipe lands; cannot be constructed. |
| 1295 | pub struct DaemonSocket { |
| 1296 | never: std::convert::Infallible, |
| 1297 | } |
| 1298 | |
| 1299 | impl DaemonSocket { |
| 1300 | #[must_use] |
| 1301 | pub fn local_path(&self) -> &Path { |
| 1302 | match self.never {} |
| 1303 | } |
| 1304 | |
| 1305 | #[must_use] |
| 1306 | pub fn shutdown_handle(&self) -> DaemonShutdownHandle { |
| 1307 | match self.never {} |
| 1308 | } |
| 1309 | |
| 1310 | pub async fn serve(self) -> Result<(), DaemonSocketError> { |
| 1311 | match self.never {} |
| 1312 | } |
| 1313 | } |
| 1314 | |
| 1315 | /// Always [`DaemonSocketError::UnsupportedPlatform`] here. |
| 1316 | pub async fn bind_daemon_socket( |
| 1317 | _options: DaemonSocketOptions, |
| 1318 | ) -> Result<DaemonSocket, DaemonSocketError> { |
| 1319 | Err(unsupported_platform()) |
| 1320 | } |
| 1321 | } |
| 1322 | |
| 1323 | #[cfg(windows)] |
| 1324 | use crate::daemon_windows as platform; |
| 1325 | pub use platform::{DaemonSocket, bind_daemon_socket}; |
| 1326 | |
| 1327 | /// `codewhale app-server --socket`: bind, announce, serve until the owner's |
| 1328 | /// `shutdown` or a termination signal. |
| 1329 | pub async fn run_daemon_socket(options: DaemonSocketOptions) -> anyhow::Result<()> { |
| 1330 | let daemon = bind_daemon_socket(options).await?; |
| 1331 | let path: &Path = daemon.local_path(); |
| 1332 | tracing::info!(path = %path.display(), "codewhale daemon listening on unix socket"); |
| 1333 | eprintln!("codewhale daemon: listening on {}", path.display()); |
| 1334 | |
| 1335 | let handle = daemon.shutdown_handle(); |
| 1336 | tokio::spawn(async move { |
| 1337 | crate::shutdown_signal().await; |
| 1338 | handle.trigger(); |
| 1339 | }); |
| 1340 | daemon.serve().await?; |
| 1341 | Ok(()) |
| 1342 | } |
| 1343 | |
| 1344 | #[cfg(any(unix, windows))] |
| 1345 | pub(crate) use platform::bind_captured_owner; |
| 1346 | |
| 1347 | #[cfg(unix)] |
| 1348 | pub async fn capture_process_start(pid: u32) -> anyhow::Result<String> { |
| 1349 | owner_work(move || codewhale_config::private_directory::unix_process_start(pid)).await |
| 1350 | } |
| 1351 | |
| 1352 | #[cfg(windows)] |
| 1353 | pub fn default_socket_path() -> Result<PathBuf, DaemonSocketError> { |
| 1354 | crate::daemon_windows::selected_pipe_path(None).map_err(DaemonSocketError::State) |
| 1355 | } |
| 1356 | |
| 1357 | #[cfg(windows)] |
| 1358 | pub async fn capture_process_start(pid: u32) -> anyhow::Result<String> { |
| 1359 | owner_work(move || { |
| 1360 | Ok( |
| 1361 | codewhale_config::windows_identity::WindowsPeerProcess::open_current_user(pid)? |
| 1362 | .start() |
| 1363 | .to_string(), |
| 1364 | ) |
| 1365 | }) |
| 1366 | .await |
| 1367 | } |
| 1368 | |
| 1369 | #[cfg(test)] |
| 1370 | mod tests { |
| 1371 | use super::*; |
| 1372 | |
| 1373 | fn inputs() -> SocketPathInputs { |
| 1374 | SocketPathInputs { |
| 1375 | explicit: None, |
| 1376 | codewhale_home_override: None, |
| 1377 | xdg_runtime_dir: None, |
| 1378 | user_home: Some(PathBuf::from("/home/whale")), |
| 1379 | macos: false, |
| 1380 | } |
| 1381 | } |
| 1382 | |
| 1383 | #[test] |
| 1384 | fn explicit_path_wins() { |
| 1385 | let resolved = resolve_socket_path(&SocketPathInputs { |
| 1386 | explicit: Some(PathBuf::from("/tmp/x.sock")), |
| 1387 | codewhale_home_override: Some(PathBuf::from("/iso")), |
| 1388 | xdg_runtime_dir: Some(PathBuf::from("/run/user/1000")), |
| 1389 | ..inputs() |
| 1390 | }) |
| 1391 | .expect("resolve"); |
| 1392 | assert_eq!(resolved, PathBuf::from("/tmp/x.sock")); |
| 1393 | } |
| 1394 | |
| 1395 | #[test] |
| 1396 | fn explicit_codewhale_home_isolates_the_daemon() { |
| 1397 | let resolved = resolve_socket_path(&SocketPathInputs { |
| 1398 | codewhale_home_override: Some(PathBuf::from("/iso/home")), |
| 1399 | xdg_runtime_dir: Some(PathBuf::from("/run/user/1000")), |
| 1400 | ..inputs() |
| 1401 | }) |
| 1402 | .expect("resolve"); |
| 1403 | assert_eq!(resolved, PathBuf::from("/iso/home/run/daemon.sock")); |
| 1404 | } |
| 1405 | |
| 1406 | #[test] |
| 1407 | fn xdg_runtime_dir_beats_home_layouts() { |
| 1408 | let resolved = resolve_socket_path(&SocketPathInputs { |
| 1409 | xdg_runtime_dir: Some(PathBuf::from("/run/user/1000")), |
| 1410 | macos: true, |
| 1411 | ..inputs() |
| 1412 | }) |
| 1413 | .expect("resolve"); |
| 1414 | assert_eq!( |
| 1415 | resolved, |
| 1416 | PathBuf::from("/run/user/1000/codewhale/daemon.sock") |
| 1417 | ); |
| 1418 | } |
| 1419 | |
| 1420 | #[test] |
| 1421 | fn macos_defaults_to_application_support() { |
| 1422 | let resolved = resolve_socket_path(&SocketPathInputs { |
| 1423 | macos: true, |
| 1424 | user_home: Some(PathBuf::from("/Users/whale")), |
| 1425 | ..inputs() |
| 1426 | }) |
| 1427 | .expect("resolve"); |
| 1428 | assert_eq!( |
| 1429 | resolved, |
| 1430 | PathBuf::from("/Users/whale/Library/Application Support/codewhale/daemon.sock") |
| 1431 | ); |
| 1432 | } |
| 1433 | |
| 1434 | #[test] |
| 1435 | fn linux_defaults_to_dot_codewhale_run() { |
| 1436 | let resolved = resolve_socket_path(&inputs()).expect("resolve"); |
| 1437 | assert_eq!( |
| 1438 | resolved, |
| 1439 | PathBuf::from("/home/whale/.codewhale/run/daemon.sock") |
| 1440 | ); |
| 1441 | } |
| 1442 | |
| 1443 | #[test] |
| 1444 | fn no_home_is_a_typed_error() { |
| 1445 | let err = resolve_socket_path(&SocketPathInputs { |
| 1446 | user_home: None, |
| 1447 | ..inputs() |
| 1448 | }) |
| 1449 | .expect_err("must fail"); |
| 1450 | assert!( |
| 1451 | matches!(err, DaemonSocketError::RuntimeDirUnavailable), |
| 1452 | "{err}" |
| 1453 | ); |
| 1454 | } |
| 1455 | |
| 1456 | #[test] |
| 1457 | fn over_long_paths_are_refused_before_bind() { |
| 1458 | let long = PathBuf::from(format!( |
| 1459 | "/{}/daemon.sock", |
| 1460 | "d".repeat(MAX_SOCKET_PATH_BYTES) |
| 1461 | )); |
| 1462 | let err = resolve_socket_path(&SocketPathInputs { |
| 1463 | explicit: Some(long.clone()), |
| 1464 | ..inputs() |
| 1465 | }) |
| 1466 | .expect_err("must fail"); |
| 1467 | match err { |
| 1468 | DaemonSocketError::PathTooLong { path, len, max } => { |
| 1469 | assert_eq!(path, long); |
| 1470 | assert!(len > max); |
| 1471 | assert_eq!(max, MAX_SOCKET_PATH_BYTES); |
| 1472 | } |
| 1473 | other => panic!("unexpected error: {other}"), |
| 1474 | } |
| 1475 | } |
| 1476 | |
| 1477 | #[test] |
| 1478 | fn unsupported_platform_error_names_the_named_pipe() { |
| 1479 | let err = unsupported_platform(); |
| 1480 | let text = err.to_string(); |
| 1481 | assert!(text.contains(WINDOWS_NAMED_PIPE), "{text}"); |
| 1482 | assert!(text.contains("not implemented"), "{text}"); |
| 1483 | } |
| 1484 | |
| 1485 | #[test] |
| 1486 | fn attach_mode_defaults_to_guest() { |
| 1487 | let params: AttachParams = |
| 1488 | serde_json::from_value(serde_json::json!({ "client": { "name": "x" } })) |
| 1489 | .expect("parse"); |
| 1490 | assert_eq!(params.mode, AttachMode::Attach); |
| 1491 | assert_eq!(params.expect_daemon_version, None); |
| 1492 | } |
| 1493 | } |
| 1494 |