Skip to main content

running_process/client/
pipe_session.rs

1//! Client-side helpers for daemon-owned pipe-backed sessions
2//! (issue #130 milestone 3).
3//!
4//! Mirrors [`crate::client::pty_session`] for the pipe case. Sessions are
5//! spawned, listed, detached, and terminated via the regular [`DaemonClient`] RPC
6//! channel. Stdin is also an RPC (`write_pipe_stdin`). Stdout/stderr are
7//! attached via [`PipeStreamAttachment`], which owns its own connection
8//! and pumps `PipeStreamFrame` payloads.
9
10use 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// ---------------------------------------------------------------------------
25// Spawn / list / terminate / write helpers
26// ---------------------------------------------------------------------------
27
28/// Request shape for spawning a daemon-owned pipe-backed session.
29#[derive(Debug, Clone)]
30pub struct PipeSpawnRequest {
31    /// Program and arguments. The first element is the executable.
32    pub argv: Vec<String>,
33    /// Working directory for the child. `None` leaves the daemon default in effect.
34    pub cwd: Option<PathBuf>,
35    /// Explicit environment variables applied after the selected base.
36    pub env: Vec<(String, String)>,
37    /// Deprecated wire-compatibility bit. Prefer [`Self::environment_policy`].
38    pub clear_inherited_env: bool,
39    /// Base environment used before applying [`Self::env`].
40    pub environment_policy: crate::EnvironmentPolicy,
41    /// Optional label used by list filters. `None` lets the daemon assign one.
42    pub originator: Option<String>,
43    /// Merge stderr into stdout instead of keeping a separately attachable stderr stream.
44    pub merge_stderr_into_stdout: bool,
45}
46
47impl PipeSpawnRequest {
48    /// Create a request that snapshots the caller environment and replaces
49    /// the daemon environment, with separate stderr.
50    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    /// Set the working directory for the child process.
63    pub fn with_cwd(mut self, cwd: impl Into<PathBuf>) -> Self {
64        self.cwd = Some(cwd.into());
65        self
66    }
67
68    /// Set the originator label stored with the session.
69    pub fn with_originator(mut self, originator: impl Into<String>) -> Self {
70        self.originator = Some(originator.into());
71        self
72    }
73
74    /// Merge stderr into stdout for this session.
75    pub fn merge_stderr(mut self) -> Self {
76        self.merge_stderr_into_stdout = true;
77        self
78    }
79
80    /// Replace the environment overlay applied when spawning the child.
81    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    /// Select the base environment for this contained remote child. `Auto`
92    /// resolves to the contained-process default, `Inherit`.
93    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/// Reply summary for a successful pipe session spawn.
107#[derive(Debug, Clone)]
108pub struct SpawnedPipeSession {
109    /// Daemon-assigned session identifier used by later pipe session RPCs.
110    pub session_id: String,
111    /// Operating-system process ID for the spawned child.
112    pub pid: u32,
113    /// Daemon-recorded creation time as seconds since the Unix epoch.
114    pub created_at: f64,
115}
116
117impl DaemonClient {
118    /// Ask the daemon to spawn a new pipe-backed child process.
119    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    /// List pipe sessions known to the daemon.
176    ///
177    /// An empty `originator_filter` returns all sessions in the current daemon scope.
178    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    /// Detach any current attachment from one pipe output stream.
205    ///
206    /// The session remains alive and the stream can be attached again later.
207    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    /// Schedule termination of a pipe session.
229    ///
230    /// The daemon accepts `0` as its default grace period before hard kill.
231    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    /// Write bytes to a session's stdin pipe.
253    ///
254    /// When `close_after` is true, the daemon closes stdin after writing `data`.
255    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
297// ---------------------------------------------------------------------------
298// PipeStreamAttachment
299// ---------------------------------------------------------------------------
300
301/// Active attachment to one stdout or stderr stream of a pipe-backed session.
302///
303/// Owns the socket after it switches into one-way streaming mode.
304pub struct PipeStreamAttachment {
305    reader: BufReader<Stream>,
306    /// Bytes from the stream backlog that the client missed before attaching.
307    pub initial_backlog: Vec<u8>,
308    /// Cumulative bytes dropped from the daemon's backlog before this attachment.
309    pub bytes_missed: u64,
310}
311
312/// Errors from opening or reading a pipe stream attachment.
313#[derive(Debug)]
314pub enum PipeAttachError {
315    /// Opening the daemon socket failed.
316    Connect(std::io::Error),
317    /// Reading or writing the length-prefixed socket stream failed.
318    Io(std::io::Error),
319    /// Decoding a daemon response or stream frame failed.
320    Decode(prost::DecodeError),
321    /// The daemon rejected the attach request.
322    Server {
323        /// Server status code returned by the daemon.
324        code: StatusCode,
325        /// Human-readable error message returned by the daemon.
326        message: String,
327    },
328    /// The daemon accepted the attach request but omitted its payload.
329    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    /// Open the scoped daemon socket and attach to a session output stream.
350    ///
351    /// When `steal` is true, the daemon evicts any existing attachment on the same stream.
352    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    /// Open an explicit daemon socket path and attach to a session output stream.
363    ///
364    /// When `steal` is true, the daemon evicts any existing attachment on the same stream.
365    pub fn attach_to(
366        socket_path: &str,
367        session_id: &str,
368        stream: PipeStreamKind,
369        steal: bool,
370    ) -> Result<Self, PipeAttachError> {
371        // Bounded connect (issue #590, cluster B): a bound-but-never-
372        // accepting daemon socket must not wedge the attaching client.
373        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        // We do not need writer after this, but keep it alive via reader's
395        // duplex socket. Drop here.
396        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    /// Block until the next stream frame arrives.
420    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    // Cap before allocating so a corrupt/desynced frame can't drive a
438    // multi-GiB allocation + unbounded read (issue #590, cluster B).
439    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;