返回 CodeWhale
daemon_windows.rs
根目录 / crates / app-server / src / daemon_windows.rs
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
369 lines RUST