1use crate::client::paths;
11use crate::client::{ClientError, DaemonClient};
12use crate::platform::ipc::Stream;
13use crate::proto::daemon::{
14 AttachPipeStreamRequest, AttachPipeStreamResponse, DaemonRequest, DaemonResponse,
15 DetachPipeStreamRequest, KeyValue, ListPipeSessionsRequest, ListPipeSessionsResponse,
16 PipeSessionInfo, PipeStreamFrame, PipeStreamKind, RequestType, SpawnPipeSessionRequest,
17 SpawnPipeSessionResponse, StatusCode, TerminatePipeSessionRequest, WritePipeStdinRequest,
18 WritePipeStdinResponse,
19};
20use prost::Message;
21use std::io::{BufReader, BufWriter, Read, Write};
22use std::path::PathBuf;
23
24#[derive(Debug, Clone)]
30pub struct PipeSpawnRequest {
31 pub argv: Vec<String>,
33 pub cwd: Option<PathBuf>,
35 pub env: Vec<(String, String)>,
37 pub clear_inherited_env: bool,
39 pub environment_policy: crate::EnvironmentPolicy,
41 pub originator: Option<String>,
43 pub merge_stderr_into_stdout: bool,
45}
46
47impl PipeSpawnRequest {
48 pub fn new<S: Into<String>>(argv: impl IntoIterator<Item = S>) -> Self {
51 Self {
52 argv: argv.into_iter().map(Into::into).collect(),
53 cwd: None,
54 env: std::env::vars().collect(),
55 clear_inherited_env: true,
56 environment_policy: crate::EnvironmentPolicy::Clear,
57 originator: None,
58 merge_stderr_into_stdout: false,
59 }
60 }
61
62 pub fn with_cwd(mut self, cwd: impl Into<PathBuf>) -> Self {
64 self.cwd = Some(cwd.into());
65 self
66 }
67
68 pub fn with_originator(mut self, originator: impl Into<String>) -> Self {
70 self.originator = Some(originator.into());
71 self
72 }
73
74 pub fn merge_stderr(mut self) -> Self {
76 self.merge_stderr_into_stdout = true;
77 self
78 }
79
80 pub fn with_envs<I, K, V>(mut self, env: I) -> Self
82 where
83 I: IntoIterator<Item = (K, V)>,
84 K: Into<String>,
85 V: Into<String>,
86 {
87 self.env = env.into_iter().map(|(k, v)| (k.into(), v.into())).collect();
88 self
89 }
90
91 pub fn with_environment_policy(mut self, policy: crate::EnvironmentPolicy) -> Self {
94 self.environment_policy = match policy {
95 crate::EnvironmentPolicy::Auto => crate::EnvironmentPolicy::Inherit,
96 explicit => explicit,
97 };
98 self.clear_inherited_env = self
99 .environment_policy
100 .legacy_clear_fallback()
101 .expect("resolved environment policy");
102 self
103 }
104}
105
106#[derive(Debug, Clone)]
108pub struct SpawnedPipeSession {
109 pub session_id: String,
111 pub pid: u32,
113 pub created_at: f64,
115}
116
117impl DaemonClient {
118 pub fn spawn_pipe_session(
120 &mut self,
121 request: &PipeSpawnRequest,
122 ) -> Result<SpawnedPipeSession, ClientError> {
123 let policy = match request.environment_policy {
124 crate::EnvironmentPolicy::Auto => crate::EnvironmentPolicy::Inherit,
125 explicit => explicit,
126 };
127 let proto = SpawnPipeSessionRequest {
128 argv: request.argv.clone(),
129 cwd: request
130 .cwd
131 .as_ref()
132 .map(|p| p.to_string_lossy().into_owned())
133 .unwrap_or_default(),
134 env: request
135 .env
136 .iter()
137 .map(|(k, v)| KeyValue {
138 key: k.clone(),
139 value: v.clone(),
140 })
141 .collect(),
142 clear_inherited_env: policy
143 .legacy_clear_fallback()
144 .map_err(|message| ClientError::Io(std::io::Error::other(message)))?,
145 originator: request.originator.clone().unwrap_or_default(),
146 merge_stderr_into_stdout: request.merge_stderr_into_stdout,
147 environment_policy: policy
148 .wire_value()
149 .map_err(|message| ClientError::Io(std::io::Error::other(message)))?,
150 };
151 let daemon_request = DaemonRequest {
152 id: self.next_request_id(),
153 r#type: RequestType::SpawnPipeSession.into(),
154 protocol_version: 1,
155 client_name: "running-process-client".into(),
156 spawn_pipe_session: Some(proto),
157 ..Default::default()
158 };
159 let response = self.send_request(daemon_request)?;
160 ensure_ok(&response)?;
161 let payload: SpawnPipeSessionResponse =
162 response
163 .spawn_pipe_session
164 .ok_or_else(|| ClientError::Server {
165 code: StatusCode::Internal,
166 message: "spawn_pipe_session response missing payload".into(),
167 })?;
168 Ok(SpawnedPipeSession {
169 session_id: payload.session_id,
170 pid: payload.pid,
171 created_at: payload.created_at,
172 })
173 }
174
175 pub fn list_pipe_sessions(
179 &mut self,
180 originator_filter: &str,
181 ) -> Result<Vec<PipeSessionInfo>, ClientError> {
182 let req = DaemonRequest {
183 id: self.next_request_id(),
184 r#type: RequestType::ListPipeSessions.into(),
185 protocol_version: 1,
186 client_name: "running-process-client".into(),
187 list_pipe_sessions: Some(ListPipeSessionsRequest {
188 originator: originator_filter.into(),
189 }),
190 ..Default::default()
191 };
192 let response = self.send_request(req)?;
193 ensure_ok(&response)?;
194 let payload: ListPipeSessionsResponse =
195 response
196 .list_pipe_sessions
197 .ok_or_else(|| ClientError::Server {
198 code: StatusCode::Internal,
199 message: "list_pipe_sessions response missing payload".into(),
200 })?;
201 Ok(payload.sessions)
202 }
203
204 pub fn detach_pipe_stream(
208 &mut self,
209 session_id: &str,
210 stream: PipeStreamKind,
211 ) -> Result<(), ClientError> {
212 let req = DaemonRequest {
213 id: self.next_request_id(),
214 r#type: RequestType::DetachPipeStream.into(),
215 protocol_version: 1,
216 client_name: "running-process-client".into(),
217 detach_pipe_stream: Some(DetachPipeStreamRequest {
218 session_id: session_id.into(),
219 stream: stream as i32,
220 }),
221 ..Default::default()
222 };
223 let response = self.send_request(req)?;
224 ensure_ok(&response)?;
225 Ok(())
226 }
227
228 pub fn terminate_pipe_session(
232 &mut self,
233 session_id: &str,
234 grace_ms: u32,
235 ) -> Result<(), ClientError> {
236 let req = DaemonRequest {
237 id: self.next_request_id(),
238 r#type: RequestType::TerminatePipeSession.into(),
239 protocol_version: 1,
240 client_name: "running-process-client".into(),
241 terminate_pipe_session: Some(TerminatePipeSessionRequest {
242 session_id: session_id.into(),
243 grace_ms,
244 }),
245 ..Default::default()
246 };
247 let response = self.send_request(req)?;
248 ensure_ok(&response)?;
249 Ok(())
250 }
251
252 pub fn write_pipe_stdin(
256 &mut self,
257 session_id: &str,
258 data: &[u8],
259 close_after: bool,
260 ) -> Result<u64, ClientError> {
261 let req = DaemonRequest {
262 id: self.next_request_id(),
263 r#type: RequestType::WritePipeStdin.into(),
264 protocol_version: 1,
265 client_name: "running-process-client".into(),
266 write_pipe_stdin: Some(WritePipeStdinRequest {
267 session_id: session_id.into(),
268 data: data.to_vec(),
269 close: close_after,
270 }),
271 ..Default::default()
272 };
273 let response = self.send_request(req)?;
274 ensure_ok(&response)?;
275 let payload: WritePipeStdinResponse =
276 response
277 .write_pipe_stdin
278 .ok_or_else(|| ClientError::Server {
279 code: StatusCode::Internal,
280 message: "write_pipe_stdin response missing payload".into(),
281 })?;
282 Ok(payload.bytes_written)
283 }
284}
285
286fn ensure_ok(response: &DaemonResponse) -> Result<(), ClientError> {
287 if response.code == StatusCode::Ok as i32 {
288 return Ok(());
289 }
290 let code = StatusCode::try_from(response.code).unwrap_or(StatusCode::UnknownRequest);
291 Err(ClientError::Server {
292 code,
293 message: response.message.clone(),
294 })
295}
296
297pub struct PipeStreamAttachment {
305 reader: BufReader<Stream>,
306 pub initial_backlog: Vec<u8>,
308 pub bytes_missed: u64,
310}
311
312#[derive(Debug)]
314pub enum PipeAttachError {
315 Connect(std::io::Error),
317 Io(std::io::Error),
319 Decode(prost::DecodeError),
321 Server {
323 code: StatusCode,
325 message: String,
327 },
328 MissingPayload,
330}
331
332impl std::fmt::Display for PipeAttachError {
333 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
334 match self {
335 Self::Connect(e) => write!(f, "pipe attach connect failed: {e}"),
336 Self::Io(e) => write!(f, "pipe attach io error: {e}"),
337 Self::Decode(e) => write!(f, "pipe attach decode error: {e}"),
338 Self::Server { code, message } => {
339 write!(f, "pipe attach server error {code:?}: {message}")
340 }
341 Self::MissingPayload => write!(f, "pipe attach response missing payload"),
342 }
343 }
344}
345
346impl std::error::Error for PipeAttachError {}
347
348impl PipeStreamAttachment {
349 pub fn attach(
353 scope_hash: Option<&str>,
354 session_id: &str,
355 stream: PipeStreamKind,
356 steal: bool,
357 ) -> Result<Self, PipeAttachError> {
358 let socket_path = paths::socket_path(scope_hash);
359 Self::attach_to(&socket_path, session_id, stream, steal)
360 }
361
362 pub fn attach_to(
366 socket_path: &str,
367 session_id: &str,
368 stream: PipeStreamKind,
369 steal: bool,
370 ) -> Result<Self, PipeAttachError> {
371 paths::make_socket_endpoint(socket_path).map_err(PipeAttachError::Connect)?;
374 let s = crate::client::deadline_io::connect_with_timeout(socket_path)
375 .map_err(PipeAttachError::Connect)?;
376 let s_clone = s.try_clone().map_err(PipeAttachError::Connect)?;
377 let mut reader = BufReader::new(s);
378 let mut writer = BufWriter::new(s_clone);
379
380 let attach_request = DaemonRequest {
381 id: 1,
382 r#type: RequestType::AttachPipeStream.into(),
383 protocol_version: 1,
384 client_name: "running-process-client".into(),
385 attach_pipe_stream: Some(AttachPipeStreamRequest {
386 session_id: session_id.into(),
387 stream: stream as i32,
388 steal,
389 }),
390 ..Default::default()
391 };
392 write_length_prefixed(&mut writer, &attach_request.encode_to_vec())
393 .map_err(PipeAttachError::Io)?;
394 drop(writer);
397
398 let response_bytes = read_length_prefixed(&mut reader).map_err(PipeAttachError::Io)?;
399 let response =
400 DaemonResponse::decode(&response_bytes[..]).map_err(PipeAttachError::Decode)?;
401 if response.code != StatusCode::Ok as i32 {
402 let code = StatusCode::try_from(response.code).unwrap_or(StatusCode::UnknownRequest);
403 return Err(PipeAttachError::Server {
404 code,
405 message: response.message,
406 });
407 }
408 let payload: AttachPipeStreamResponse = response
409 .attach_pipe_stream
410 .ok_or(PipeAttachError::MissingPayload)?;
411
412 Ok(Self {
413 reader,
414 initial_backlog: payload.backlog,
415 bytes_missed: payload.bytes_missed,
416 })
417 }
418
419 pub fn recv_frame(&mut self) -> Result<PipeStreamFrame, PipeAttachError> {
421 let bytes = read_length_prefixed(&mut self.reader).map_err(PipeAttachError::Io)?;
422 PipeStreamFrame::decode(&bytes[..]).map_err(PipeAttachError::Decode)
423 }
424}
425
426fn write_length_prefixed<W: Write>(w: &mut W, payload: &[u8]) -> Result<(), std::io::Error> {
427 let len = payload.len() as u32;
428 w.write_all(&len.to_be_bytes())?;
429 w.write_all(payload)?;
430 w.flush()
431}
432
433fn read_length_prefixed<R: Read>(r: &mut R) -> Result<Vec<u8>, std::io::Error> {
434 let mut len_buf = [0u8; 4];
435 r.read_exact(&mut len_buf)?;
436 let len = u32::from_be_bytes(len_buf) as usize;
437 crate::client::deadline_io::check_frame_len(len)?;
440 let mut buf = vec![0u8; len];
441 r.read_exact(&mut buf)?;
442 Ok(buf)
443}
444
445#[cfg(test)]
446mod tests {
447 use super::*;
448
449 #[test]
450 fn pipe_spawn_request_defaults_to_client_snapshot_replacement() {
451 let request = PipeSpawnRequest::new(["echo", "hi"]);
452 assert_eq!(request.environment_policy, crate::EnvironmentPolicy::Clear);
453 assert!(request.clear_inherited_env);
454 assert!(!request.env.is_empty());
455 }
456
457 #[test]
458 fn pipe_spawn_request_dual_writes_explicit_policy() {
459 let inherit = PipeSpawnRequest::new(["echo"])
460 .with_environment_policy(crate::EnvironmentPolicy::Inherit);
461 assert_eq!(
462 inherit.environment_policy,
463 crate::EnvironmentPolicy::Inherit
464 );
465 assert!(!inherit.clear_inherited_env);
466
467 let baseline = PipeSpawnRequest::new(["echo"])
468 .with_environment_policy(crate::EnvironmentPolicy::UserBaseline);
469 assert_eq!(
470 baseline.environment_policy,
471 crate::EnvironmentPolicy::UserBaseline
472 );
473 assert!(baseline.clear_inherited_env);
474 }
475}
476
477#[cfg(test)]
478#[path = "../tests/client_pipe_session_coverage.rs"]
479mod coverage_tests;