| 1 | //! Windows kernel pipe frontend for the existing authenticated control loop. |
| 2 | //! Pipe names route; private receipts plus held kernel process/token identity |
| 3 | //! select the actual owner. No bearer or Engine authority crosses this wire. |
| 4 | |
| 5 | use crate::daemon_socket::{ |
| 6 | AuthorizedPeer, ConnectionTasks, DaemonShutdownHandle, DaemonSocketError, DaemonSocketOptions, |
| 7 | SocketPathInputs, connection_context, owner_work, run_authorized_connection, |
| 8 | }; |
| 9 | use crate::{AppState, AppTransport, build_state_off_runtime}; |
| 10 | use anyhow::{Context as _, Result}; |
| 11 | use codewhale_config::private_directory::PrivateDirectory; |
| 12 | use codewhale_config::windows_identity::{CurrentWindowsUser, OwnerOnlyAcl, WindowsPeerProcess}; |
| 13 | use codewhale_protocol::RuntimeOwnerReceipt; |
| 14 | use std::fs::File; |
| 15 | use std::os::windows::io::AsRawHandle as _; |
| 16 | use std::path::{Path, PathBuf}; |
| 17 | use std::sync::Arc; |
| 18 | use tokio::net::windows::named_pipe::{ClientOptions, NamedPipeServer, ServerOptions}; |
| 19 | use tokio::sync::watch; |
| 20 | |
| 21 | pub(crate) fn owner_directory() -> Result<PathBuf> { |
| 22 | let input = SocketPathInputs::from_environment(None)?; |
| 23 | let home = input |
| 24 | .codewhale_home_override |
| 25 | .or_else(|| { |
| 26 | input |
| 27 | .user_home |
| 28 | .map(|p| p.join(codewhale_paths::CODEWHALE_APP_DIR)) |
| 29 | }) |
| 30 | .context("selected Windows owner home unavailable")?; |
| 31 | anyhow::ensure!( |
| 32 | home.is_absolute(), |
| 33 | "selected Windows owner home must be absolute" |
| 34 | ); |
| 35 | Ok(home.join("run")) |
| 36 | } |
| 37 | |
| 38 | pub(crate) fn selected_pipe_path(explicit: Option<PathBuf>) -> Result<PathBuf> { |
| 39 | if let Some(path) = explicit { |
| 40 | let text = path.to_str().context("invalid Windows pipe name")?; |
| 41 | anyhow::ensure!( |
| 42 | text.starts_with(r"\\.\pipe\") |
| 43 | && text.len() < 240 |
| 44 | && !text[9..].contains(['\\', '/', '\0']), |
| 45 | "control pipe must be one local named-pipe component" |
| 46 | ); |
| 47 | return Ok(path); |
| 48 | } |
| 49 | use sha2::{Digest, Sha256}; |
| 50 | let home = owner_directory()?; |
| 51 | let principal = CurrentWindowsUser::open()?.sid_string()?; |
| 52 | let mut digest = Sha256::new(); |
| 53 | digest.update(b"Codewhale/owner-pipe/v1\0"); |
| 54 | digest.update(home.as_os_str().as_encoded_bytes()); |
| 55 | digest.update(b"\0"); |
| 56 | digest.update(principal.as_bytes()); |
| 57 | let suffix: String = digest |
| 58 | .finalize() |
| 59 | .iter() |
| 60 | .map(|byte| format!("{byte:02x}")) |
| 61 | .collect(); |
| 62 | Ok(PathBuf::from(format!(r"\\.\pipe\codewhale-owner-{suffix}"))) |
| 63 | } |
| 64 | |
| 65 | fn create_pipe(path: &Path, first: bool) -> Result<NamedPipeServer> { |
| 66 | use windows_sys::Win32::Storage::FileSystem::FILE_ALL_ACCESS; |
| 67 | let acl = OwnerOnlyAcl::new(FILE_ALL_ACCESS, false)?; |
| 68 | let mut options = ServerOptions::new(); |
| 69 | options |
| 70 | .first_pipe_instance(first) |
| 71 | .reject_remote_clients(true) |
| 72 | .max_instances(65) |
| 73 | .in_buffer_size(16384) |
| 74 | .out_buffer_size(16384); |
| 75 | // Security is supplied at creation. There is never a default-DACL window. |
| 76 | acl.with_security_attributes(|attributes| unsafe { |
| 77 | options |
| 78 | .create_with_security_attributes_raw(path, attributes) |
| 79 | .map_err(Into::into) |
| 80 | }) |
| 81 | } |
| 82 | |
| 83 | struct ReceiptGuard { |
| 84 | parent: Arc<PrivateDirectory>, |
| 85 | file: Option<File>, |
| 86 | } |
| 87 | impl ReceiptGuard { |
| 88 | async fn retire(&mut self) -> Result<()> { |
| 89 | if let Some(file) = self.file.take() { |
| 90 | let parent = self.parent.clone(); |
| 91 | tokio::spawn(async move { |
| 92 | owner_work(move || { |
| 93 | parent.retire_private_receipt("daemon.owner.json", &file)?; |
| 94 | Ok(()) |
| 95 | }) |
| 96 | .await |
| 97 | }) |
| 98 | .await??; |
| 99 | } |
| 100 | Ok(()) |
| 101 | } |
| 102 | } |
| 103 | impl Drop for ReceiptGuard { |
| 104 | fn drop(&mut self) { |
| 105 | let Some(file) = self.file.take() else { return }; |
| 106 | let parent = self.parent.clone(); |
| 107 | let retire = move || -> Result<()> { |
| 108 | parent.retire_private_receipt("daemon.owner.json", &file)?; |
| 109 | Ok(()) |
| 110 | }; |
| 111 | if let Ok(runtime) = tokio::runtime::Handle::try_current() { |
| 112 | runtime.spawn(async move { |
| 113 | if let Err(error) = owner_work(retire).await { |
| 114 | tracing::warn!(%error,"private Windows owner retirement uncertain"); |
| 115 | } |
| 116 | }); |
| 117 | } else if let Err(error) = retire() { |
| 118 | tracing::warn!(%error,"private Windows owner retirement uncertain"); |
| 119 | } |
| 120 | } |
| 121 | } |
| 122 | |
| 123 | pub struct DaemonSocket { |
| 124 | pipe: NamedPipeServer, |
| 125 | path: PathBuf, |
| 126 | state: AppState, |
| 127 | shutdown: Arc<watch::Sender<bool>>, |
| 128 | owner: Option<RuntimeOwnerReceipt>, |
| 129 | receipt: Option<ReceiptGuard>, |
| 130 | } |
| 131 | impl DaemonSocket { |
| 132 | pub fn local_path(&self) -> &Path { |
| 133 | &self.path |
| 134 | } |
| 135 | pub fn shutdown_handle(&self) -> DaemonShutdownHandle { |
| 136 | DaemonShutdownHandle(self.shutdown.clone()) |
| 137 | } |
| 138 | pub async fn serve(self) -> Result<(), DaemonSocketError> { |
| 139 | let Self { |
| 140 | mut pipe, |
| 141 | path, |
| 142 | state, |
| 143 | shutdown, |
| 144 | owner, |
| 145 | mut receipt, |
| 146 | } = self; |
| 147 | let context = connection_context( |
| 148 | state, |
| 149 | path.clone(), |
| 150 | DaemonShutdownHandle(shutdown.clone()), |
| 151 | owner, |
| 152 | ) |
| 153 | .await |
| 154 | .map_err(DaemonSocketError::State)?; |
| 155 | let mut stopping = shutdown.subscribe(); |
| 156 | let mut connections = ConnectionTasks::default(); |
| 157 | loop { |
| 158 | if *stopping.borrow() { |
| 159 | break; |
| 160 | } |
| 161 | tokio::select! { |
| 162 | result=pipe.connect()=>{ |
| 163 | result.map_err(|error|DaemonSocketError::State(error.into()))?; |
| 164 | if connections.at_capacity() {pipe.disconnect().map_err(|error|DaemonSocketError::State(error.into()))?;continue;} |
| 165 | // Retain the connected instance until the successor exists; |
| 166 | // the named endpoint is never released between accepts. |
| 167 | let successor_path=path.clone(); |
| 168 | let successor=owner_work(move ||create_pipe(&successor_path,false)).await.map_err(DaemonSocketError::State)?; |
| 169 | let accepted=std::mem::replace(&mut pipe,successor); |
| 170 | let context=context.clone(); |
| 171 | connections.spawn(async move { |
| 172 | use windows_sys::Win32::System::Pipes::GetNamedPipeClientProcessId; |
| 173 | let mut pid=0; |
| 174 | if unsafe {GetNamedPipeClientProcessId(accepted.as_raw_handle(),&mut pid)}==0 {return} |
| 175 | let process=match owner_work(move ||WindowsPeerProcess::open_current_user(pid)).await {Ok(process)=>process,Err(_)=>return}; |
| 176 | let start=process.start().to_string();let (read,write)=tokio::io::split(accepted); |
| 177 | run_authorized_connection(context,read,write,AuthorizedPeer {pid,start,process}).await; |
| 178 | }); |
| 179 | } |
| 180 | changed=stopping.changed()=>{if changed.is_err() || *stopping.borrow() {break}} |
| 181 | } |
| 182 | } |
| 183 | connections.shutdown().await; |
| 184 | drop(pipe); |
| 185 | if let Some(receipt) = receipt.as_mut() { |
| 186 | receipt.retire().await.map_err(DaemonSocketError::State)?; |
| 187 | } |
| 188 | Ok(()) |
| 189 | } |
| 190 | } |
| 191 | |
| 192 | pub async fn bind_daemon_socket( |
| 193 | options: DaemonSocketOptions, |
| 194 | ) -> Result<DaemonSocket, DaemonSocketError> { |
| 195 | bind(options, None).await |
| 196 | } |
| 197 | pub(crate) async fn bind_captured_owner( |
| 198 | state: AppState, |
| 199 | owner: RuntimeOwnerReceipt, |
| 200 | ) -> Result<DaemonSocket, DaemonSocketError> { |
| 201 | bind( |
| 202 | DaemonSocketOptions { |
| 203 | socket_path: Some(owner.socket_path.clone()), |
| 204 | config_path: None, |
| 205 | }, |
| 206 | Some((state, owner)), |
| 207 | ) |
| 208 | .await |
| 209 | } |
| 210 | async fn bind( |
| 211 | options: DaemonSocketOptions, |
| 212 | captured: Option<(AppState, RuntimeOwnerReceipt)>, |
| 213 | ) -> Result<DaemonSocket, DaemonSocketError> { |
| 214 | let explicit = options.socket_path; |
| 215 | let path = owner_work(move || selected_pipe_path(explicit)) |
| 216 | .await |
| 217 | .map_err(DaemonSocketError::State)?; |
| 218 | let selected = path.clone(); |
| 219 | let pipe = owner_work(move || create_pipe(&selected, true)) |
| 220 | .await |
| 221 | .map_err(DaemonSocketError::State)?; |
| 222 | let (state, owner) = match captured { |
| 223 | Some((state, owner)) => (state, Some(owner)), |
| 224 | None => ( |
| 225 | build_state_off_runtime(options.config_path, None, AppTransport::Socket) |
| 226 | .await |
| 227 | .map_err(DaemonSocketError::State)?, |
| 228 | None, |
| 229 | ), |
| 230 | }; |
| 231 | let receipt = if let Some(owner) = owner.as_ref() { |
| 232 | let owner = owner.clone(); |
| 233 | Some( |
| 234 | owner_work(move || { |
| 235 | let parent = Arc::new(PrivateDirectory::admit(&owner_directory()?)?); |
| 236 | let mut guard = ReceiptGuard { |
| 237 | parent: parent.clone(), |
| 238 | file: None, |
| 239 | }; |
| 240 | anyhow::ensure!( |
| 241 | parent.is_at_selected_path()?, |
| 242 | "selected Windows owner parent changed" |
| 243 | ); |
| 244 | if let Some((bytes, file)) = |
| 245 | parent.read_private_receipt("daemon.owner.json", 16384)? |
| 246 | { |
| 247 | let previous: RuntimeOwnerReceipt = serde_json::from_slice(&bytes)?; |
| 248 | anyhow::ensure!( |
| 249 | previous.version == owner.version |
| 250 | && previous.data_dir == owner.data_dir |
| 251 | && previous.execution_scope == owner.execution_scope |
| 252 | && previous.socket_path == owner.socket_path, |
| 253 | "old Windows owner receipt belongs to another selected store" |
| 254 | ); |
| 255 | anyhow::ensure!( |
| 256 | parent.retire_private_receipt("daemon.owner.json", &file)?, |
| 257 | "Windows owner receipt changed" |
| 258 | ); |
| 259 | drop(file); |
| 260 | } |
| 261 | let bytes = serde_json::to_vec(&owner)?; |
| 262 | anyhow::ensure!(bytes.len() <= 16384, "Windows owner receipt exceeds bound"); |
| 263 | parent.write_owned_file("daemon.owner.json", &bytes, false)?; |
| 264 | guard.file = Some( |
| 265 | parent |
| 266 | .read_private_receipt("daemon.owner.json", 16384)? |
| 267 | .context("Windows owner publication unavailable")? |
| 268 | .1, |
| 269 | ); |
| 270 | anyhow::ensure!( |
| 271 | parent.is_at_selected_path()?, |
| 272 | "selected Windows owner parent changed during publication" |
| 273 | ); |
| 274 | Ok(guard) |
| 275 | }) |
| 276 | .await |
| 277 | .map_err(DaemonSocketError::State)?, |
| 278 | ) |
| 279 | } else { |
| 280 | None |
| 281 | }; |
| 282 | let (shutdown, _) = watch::channel(false); |
| 283 | Ok(DaemonSocket { |
| 284 | pipe, |
| 285 | path, |
| 286 | state, |
| 287 | shutdown: Arc::new(shutdown), |
| 288 | owner, |
| 289 | receipt, |
| 290 | }) |
| 291 | } |
| 292 | |
| 293 | pub(crate) async fn connect_owner( |
| 294 | owner: &RuntimeOwnerReceipt, |
| 295 | ) -> Result<( |
| 296 | tokio::net::windows::named_pipe::NamedPipeClient, |
| 297 | WindowsPeerProcess, |
| 298 | )> { |
| 299 | let path = owner.socket_path.clone(); |
| 300 | let expected = owner.clone(); |
| 301 | owner_work(move || { |
| 302 | let client = ClientOptions::new() |
| 303 | .open(&path) |
| 304 | .context("opening selected Windows owner; refusing fallback")?; |
| 305 | use windows_sys::Win32::System::Pipes::GetNamedPipeServerProcessId; |
| 306 | let mut pid = 0; |
| 307 | anyhow::ensure!( |
| 308 | unsafe { GetNamedPipeServerProcessId(client.as_raw_handle(), &mut pid) } != 0 |
| 309 | && pid == expected.pid, |
| 310 | "connected Windows server is not the selected owner" |
| 311 | ); |
| 312 | let process = WindowsPeerProcess::open_current_user(pid)?; |
| 313 | anyhow::ensure!( |
| 314 | process.start() == expected.process_start && process.principal()? == expected.principal, |
| 315 | "selected Windows server generation/principal changed" |
| 316 | ); |
| 317 | Ok((client, process)) |
| 318 | }) |
| 319 | .await |
| 320 | } |
| 321 | |
| 322 | #[cfg(test)] |
| 323 | mod tests { |
| 324 | use super::*; |
| 325 | |
| 326 | #[tokio::test] |
| 327 | async fn first_instance_pipe_refuses_competing_listener_and_reopens_after_close() { |
| 328 | let path = PathBuf::from(format!( |
| 329 | r"\\.\pipe\codewhale-owner-test-{}", |
| 330 | uuid::Uuid::new_v4() |
| 331 | )); |
| 332 | let first = create_pipe(&path, true).expect("current-user protected first pipe"); |
| 333 | assert!( |
| 334 | create_pipe(&path, true).is_err(), |
| 335 | "never join/substitute another selected owner instance" |
| 336 | ); |
| 337 | drop(first); |
| 338 | let reopened = |
| 339 | create_pipe(&path, true).expect("closed instance releases exact kernel endpoint"); |
| 340 | drop(reopened); |
| 341 | } |
| 342 | |
| 343 | #[test] |
| 344 | fn local_pipe_selector_refuses_remote_and_nested_names() { |
| 345 | for name in [ |
| 346 | r"\\remote\pipe\codewhale", |
| 347 | r"\\.\pipe\nested\owner", |
| 348 | r"\\.\pipe\owner/name", |
| 349 | ] { |
| 350 | assert!(selected_pipe_path(Some(PathBuf::from(name))).is_err()); |
| 351 | } |
| 352 | } |
| 353 | |
| 354 | #[test] |
| 355 | fn held_peer_captures_actual_current_user_process_generation() { |
| 356 | assert!(WindowsPeerProcess::open_current_user(0).is_err()); |
| 357 | let peer = |
| 358 | WindowsPeerProcess::open_current_user(std::process::id()).expect("held process/token"); |
| 359 | assert_eq!(peer.pid(), std::process::id()); |
| 360 | assert!(peer.start().starts_with("windows:")); |
| 361 | assert_eq!( |
| 362 | peer.principal().unwrap(), |
| 363 | CurrentWindowsUser::open().unwrap().sid_string().unwrap() |
| 364 | ); |
| 365 | peer.check_current_user() |
| 366 | .expect("held generation remains current"); |
| 367 | } |
| 368 | } |
| 369 |